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