init-draft

This commit is contained in:
Ubuntu
2020-07-09 01:29:40 +00:00
parent d6b6842ceb
commit 5769a36640
5 changed files with 257 additions and 123 deletions
+59 -59
View File
@@ -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: "<pre>{{ toPrettyJson . }}</pre>"
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
# - name: "alert"
# opsgenie:
# apiKey: ""
# priority: "P3"
# message: "Event {{ .Reason }} for {{ .InvolvedObject.Namespace }}/{{ .InvolvedObject.Name }} on K8s cluster"
# alias: "{{ .UID }}"
# description: "<pre>{{ toPrettyJson . }}</pre>"
# 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
+7
View File
@@ -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.")
}
+4 -2
View File
@@ -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.
+185 -62
View File
@@ -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() {
+2
View File
@@ -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