From 5769a3664097866e87f1d48ec6aa49baadfc9674 Mon Sep 17 00:00:00 2001 From: Ubuntu Date: Thu, 9 Jul 2020 01:16:23 +0000 Subject: [PATCH] init-draft --- config.yaml | 118 +++++++++--------- pkg/batch/writer.go | 7 ++ pkg/kube/event.go | 6 +- pkg/sinks/elasticsearch.go | 247 +++++++++++++++++++++++++++---------- pkg/sinks/reciever.go | 2 + 5 files changed, 257 insertions(+), 123 deletions(-) diff --git a/config.yaml b/config.yaml index ca7a6ba..ca2f83e 100644 --- a/config.yaml +++ b/config.yaml @@ -25,62 +25,62 @@ receivers: hosts: - "http://localhost:9200" indexFormat: "kube-events-{2006-01-02}" - - name: "alert" - opsgenie: - apiKey: "" - priority: "P3" - message: "Event {{ .Reason }} for {{ .InvolvedObject.Namespace }}/{{ .InvolvedObject.Name }} on K8s cluster" - alias: "{{ .UID }}" - description: "
{{ toPrettyJson . }}
" - tags: - - "event" - - "{{ .Reason }}" - - "{{ .InvolvedObject.Kind }}" - - "{{ .InvolvedObject.Name }}" - - name: "slack" - slack: - token: "" - channel: "#mustafa-test" - message: "Received a Kubernetes Event {{ .Message}}" - fields: - message: "{{ .Message }}" - namespace: "{{ .Namespace }}" - reason: "{{ .Reason }}" - object: "{{ .Namespace }}" - - name: "pipe" - webhook: - endpoint: "http://localhost:3000" - headers: - X-API-KEY: "123-456-OPSGENIE-789-ABC" - User-Agent: "kube-event-exporter 1.0" - streamName: "applicationMetric" - layout: - endpoint: "localhost2" - eventType: "kube-event" - createdAt: "{{ .GetTimestampMs }}" - details: - message: "{{ .Message }}" - reason: "{{ .Reason }}" - tip: "{{ .Type }}" - count: "{{ .Count }}" - kind: "{{ .InvolvedObject.Kind }}" - name: "{{ .InvolvedObject.Name }}" - namespace: "{{ .Namespace }}" - component: "{{ .Source.Component }}" - host: "{{ .Source.Host }}" - labels: "{{ toJson .InvolvedObject.Labels}}" - - name: "kafka" - kafka: - topic: "kube-event" - brokers: - - "localhost:9092" - tls: - enable: false - certFile: "kafka-client.crt" - keyFile: "kafka-client.key" - caFile: "kafka-ca.crt" - - name: "pubsub" - pubsub: - gcloud_project_id: "my-project" - topic: "kube-event" - create_topic: False \ No newline at end of file +# - name: "alert" +# opsgenie: +# apiKey: "" +# priority: "P3" +# message: "Event {{ .Reason }} for {{ .InvolvedObject.Namespace }}/{{ .InvolvedObject.Name }} on K8s cluster" +# alias: "{{ .UID }}" +# description: "
{{ toPrettyJson . }}
" +# tags: +# - "event" +# - "{{ .Reason }}" +# - "{{ .InvolvedObject.Kind }}" +# - "{{ .InvolvedObject.Name }}" +# - name: "slack" +# slack: +# token: "" +# channel: "#mustafa-test" +# message: "Received a Kubernetes Event {{ .Message}}" +# fields: +# message: "{{ .Message }}" +# namespace: "{{ .Namespace }}" +# reason: "{{ .Reason }}" +# object: "{{ .Namespace }}" +# - name: "pipe" +# webhook: +# endpoint: "http://localhost:3000" +# headers: +# X-API-KEY: "123-456-OPSGENIE-789-ABC" +# User-Agent: "kube-event-exporter 1.0" +# streamName: "applicationMetric" +# layout: +# endpoint: "localhost2" +# eventType: "kube-event" +# createdAt: "{{ .GetTimestampMs }}" +# details: +# message: "{{ .Message }}" +# reason: "{{ .Reason }}" +# tip: "{{ .Type }}" +# count: "{{ .Count }}" +# kind: "{{ .InvolvedObject.Kind }}" +# name: "{{ .InvolvedObject.Name }}" +# namespace: "{{ .Namespace }}" +# component: "{{ .Source.Component }}" +# host: "{{ .Source.Host }}" +# labels: "{{ toJson .InvolvedObject.Labels}}" +# - name: "kafka" +# kafka: +# topic: "kube-event" +# brokers: +# - "localhost:9092" +# tls: +# enable: false +# certFile: "kafka-client.crt" +# keyFile: "kafka-client.key" +# caFile: "kafka-ca.crt" +# - name: "pubsub" +# pubsub: +# gcloud_project_id: "my-project" +# topic: "kube-event" +# create_topic: False diff --git a/pkg/batch/writer.go b/pkg/batch/writer.go index e484d33..d6b8199 100644 --- a/pkg/batch/writer.go +++ b/pkg/batch/writer.go @@ -3,6 +3,7 @@ package batch import ( "context" "time" + "github.com/rs/zerolog/log" ) @@ -54,6 +55,7 @@ func (w *Writer) Start() { for shouldGoOn { select { case item := <-w.items: + log.Info().Msgf("w.len: %v, w.cfg.BatchSize: %v", w.len, w.cfg.BatchSize) if w.len >= w.cfg.BatchSize { w.processBuffer(context.Background()) w.len = 0 @@ -67,6 +69,7 @@ func (w *Writer) Start() { w.stopDone <- true ticker.Stop() case <-ticker.C: + log.Info().Msgf("ticker") w.processBuffer(context.Background()) } } @@ -78,6 +81,8 @@ func (w *Writer) processBuffer(ctx context.Context) { return } + log.Info().Msgf("processBuffer, w.len: %v", w.len) + // Need to copy the underlying item to another slice slice := make([]interface{}, w.len) for i := 0; i < w.len; i++ { @@ -107,6 +112,7 @@ func (w *Writer) processBuffer(ctx context.Context) { } } + log.Info().Msgf("processBuffer, newItemsCount: %v", newItemsCount) w.len = newItemsCount // TODO(makin) an edge case, if all items fail, and the buffer is full, new item cannot be added to buffer. } @@ -122,4 +128,5 @@ func (w *Writer) Submit(items ...interface{}) { for _, item := range items { w.items <- item } + log.Info().Msgf("Submit done.") } diff --git a/pkg/kube/event.go b/pkg/kube/event.go index e74debd..660f53c 100644 --- a/pkg/kube/event.go +++ b/pkg/kube/event.go @@ -11,10 +11,12 @@ type EnhancedEvent struct { InvolvedObject EnhancedObjectReference `json:"involvedObject"` } + +// TODO(vsbus): restore this part but for BQ sink find a way to convert map into list of KV pairs before marshaling to JSON. type EnhancedObjectReference struct { corev1.ObjectReference `json:",inline"` - Labels map[string]string `json:"labels,omitempty"` - Annotations map[string]string `json:"annotations,omitempty"` + Labels map[string]string `json:"-"` + Annotations map[string]string `json:"-"` } // ToJSON does not return an error because we are %99 confident it is JSON serializable. diff --git a/pkg/sinks/elasticsearch.go b/pkg/sinks/elasticsearch.go index 482420f..dd5b17d 100644 --- a/pkg/sinks/elasticsearch.go +++ b/pkg/sinks/elasticsearch.go @@ -1,15 +1,22 @@ package sinks import ( - "bytes" + // "bytes" + "bufio" + "os" + "cloud.google.com/go/bigquery" "context" - "crypto/tls" - "encoding/json" + "encoding/json" + // "crypto/tls" + // "encoding/json" + "fmt" "github.com/elastic/go-elasticsearch/v7" - "github.com/elastic/go-elasticsearch/v7/esapi" + // "github.com/elastic/go-elasticsearch/v7/esapi" "github.com/opsgenie/kubernetes-event-exporter/pkg/kube" - "io/ioutil" - "net/http" + "github.com/opsgenie/kubernetes-event-exporter/pkg/batch" + "github.com/rs/zerolog/log" + // "io/ioutil" + // "net/http" "regexp" "strings" "time" @@ -33,46 +40,131 @@ type ElasticsearchConfig struct { } `yaml:"tls"` } +func writeBatchToJsonFile(path string, items []interface{}) error { + file, err := os.Create(path) + if err != nil { + return err + } + defer file.Close() + + writer := bufio.NewWriter(file) + for i := 0; i < len(items); i++ { + jsonBytes, err := json.Marshal(items[i]) + if err != nil { + log.Warn().Msgf("Failed to convert item to json: %v", items[i]) + } else { + fmt.Fprintln(writer, string(jsonBytes)) + } + } + return writer.Flush(); +} + +func importJsonFromFile(filename string) error { + projectID := "foo" + datasetID := "bar" + tableID := "baz" + ctx := context.Background() + client, err := bigquery.NewClient(ctx, projectID) + if err != nil { + return fmt.Errorf("bigquery.NewClient: %v", err) + } + defer client.Close() + + f, err := os.Open(filename) + if err != nil { + return err + } + source := bigquery.NewReaderSource(f) + source.SourceFormat = bigquery.JSON + source.AutoDetect = true // Allow BigQuery to determine schema. + + loader := client.Dataset(datasetID).Table(tableID).LoaderFrom(source) + + log.Info().Msgf("loader.Run...") + job, err := loader.Run(ctx) + if err != nil { + return err + } + log.Info().Msgf("loader.Wait...") + status, err := job.Wait(ctx) + if err != nil { + return err + } + log.Info().Msgf("loader done.") + if err := status.Err(); err != nil { + return err + } + return nil +} + + func NewElasticsearch(cfg *ElasticsearchConfig) (*Elasticsearch, error) { - var caCert []byte + log.Info().Msgf("NewElasticsearch cfg: %v", cfg) + // var caCert []byte - if len(cfg.TLS.CaFile) > 0 { - readFile, err := ioutil.ReadFile(cfg.TLS.CaFile) - if err != nil { - return nil, err - } - caCert = readFile + // if len(cfg.TLS.CaFile) > 0 { + // readFile, err := ioutil.ReadFile(cfg.TLS.CaFile) + // if err != nil { + // return nil, err + // } + // caCert = readFile + // } + + // tlsClientConfig := &tls.Config{ + // InsecureSkipVerify: cfg.TLS.InsecureSkipVerify, + // ServerName: cfg.TLS.ServerName, + // } + // tlsClientConfig.RootCAs.AppendCertsFromPEM(caCert) + + // client, err := elasticsearch.NewClient(elasticsearch.Config{ + // Addresses: cfg.Hosts, + // Username: cfg.Username, + // Password: cfg.Password, + // CloudID: cfg.CloudID, + // APIKey: cfg.APIKey, + // Transport: &http.Transport{ + // TLSClientConfig: tlsClientConfig, + // }, + // }) + // if err != nil { + // return nil, err + // } + + myfunc := func(ctx context.Context, items []interface{}) []bool { + res := make([]bool, len(items)) + for i := 0; i < len(items); i++ { + res[i] = true + } + path := "/tmp/batch.json" + if err := writeBatchToJsonFile(path, items); err != nil { + log.Error().Msgf("Failed to write JSON file: %v", err) + } + if err := importJsonFromFile(path); err != nil { + log.Error().Msgf("BigQuery load failed: %v", err) + } + return res } - - tlsClientConfig := &tls.Config{ - InsecureSkipVerify: cfg.TLS.InsecureSkipVerify, - ServerName: cfg.TLS.ServerName, - } - tlsClientConfig.RootCAs.AppendCertsFromPEM(caCert) - - client, err := elasticsearch.NewClient(elasticsearch.Config{ - Addresses: cfg.Hosts, - Username: cfg.Username, - Password: cfg.Password, - CloudID: cfg.CloudID, - APIKey: cfg.APIKey, - Transport: &http.Transport{ - TLSClientConfig: tlsClientConfig, - }, - }) - if err != nil { - return nil, err - } - + batchWriter := batch.NewWriter( + batch.WriterConfig{ + BatchSize: 1000, + MaxRetries: 3, + Interval: time.Duration(10) * time.Second, + Timeout: time.Duration(60) * time.Second, + }, + myfunc, + ) + batchWriter.Start() return &Elasticsearch{ - client: client, - cfg: cfg, + client: nil, + cfg: nil, + batchWriter: batchWriter, }, nil } type Elasticsearch struct { client *elasticsearch.Client cfg *ElasticsearchConfig + batchWriter *batch.Writer } var regex = regexp.MustCompile(`(?s){(.*)}`) @@ -95,37 +187,68 @@ func formatIndexName(pattern string, when time.Time) string { return builder.String() } + +func importJSONAutodetectSchema(projectID, datasetID, tableID string) error { + // projectID := "my-project-id" + // datasetID := "mydataset" + // tableID := "mytable" + ctx := context.Background() + client, err := bigquery.NewClient(ctx, projectID) + if err != nil { + return fmt.Errorf("bigquery.NewClient: %v", err) + } + defer client.Close() + + gcsRef := bigquery.NewGCSReference("gs://cloud-samples-data/bigquery/us-states/us-states.json") + gcsRef.SourceFormat = bigquery.JSON + gcsRef.AutoDetect = true + loader := client.Dataset(datasetID).Table(tableID).LoaderFrom(gcsRef) + loader.WriteDisposition = bigquery.WriteEmpty + + job, err := loader.Run(ctx) + if err != nil { + return err + } + status, err := job.Wait(ctx) + if err != nil { + return err + } + + if status.Err() != nil { + return fmt.Errorf("job completed with error: %v", status.Err()) + } + return nil +} + func (e *Elasticsearch) Send(ctx context.Context, ev *kube.EnhancedEvent) error { - b, err := json.Marshal(ev) - if err != nil { - return err - } + log.Info().Msgf("add to buffer...") + e.batchWriter.Submit(ev) + return nil + // var index string + // if len(e.cfg.IndexFormat) > 0 { + // now := time.Now() + // index = formatIndexName(e.cfg.IndexFormat, now) + // } else { + // index = e.cfg.Index + // } - var index string - if len(e.cfg.IndexFormat) > 0 { - now := time.Now() - index = formatIndexName(e.cfg.IndexFormat, now) - } else { - index = e.cfg.Index - } + // req := esapi.IndexRequest{ + // Body: bytes.NewBuffer(b), + // Index: index, + // } - req := esapi.IndexRequest{ - Body: bytes.NewBuffer(b), - Index: index, - } + // if e.cfg.UseEventID { + // req.DocumentID = string(ev.UID) + // } - if e.cfg.UseEventID { - req.DocumentID = string(ev.UID) - } + // resp, err := req.Do(ctx, e.client) + // if err != nil { + // return err + // } - resp, err := req.Do(ctx, e.client) - if err != nil { - return err - } - - defer resp.Body.Close() - _ = resp.Body - return nil + // defer resp.Body.Close() + // _ = resp.Body + // return nil } func (e *Elasticsearch) Close() { diff --git a/pkg/sinks/reciever.go b/pkg/sinks/reciever.go index 72a9417..239878b 100644 --- a/pkg/sinks/reciever.go +++ b/pkg/sinks/reciever.go @@ -1,6 +1,7 @@ package sinks import "errors" +import "fmt" // Receiver allows receiving type ReceiverConfig struct { @@ -24,6 +25,7 @@ func (r *ReceiverConfig) Validate() error { } func (r *ReceiverConfig) GetSink() (Sink, error) { + fmt.Println("GetSink ", r); if r.InMemory != nil { // This reference is used for test purposes to count the events in the sink. // It should not be used in production since it will only cause memory leak and (b)OOM