/
data_recorder_kinesis.go
81 lines (70 loc) · 2.38 KB
/
data_recorder_kinesis.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
package handler
import (
producer "github.com/a8m/kinesis-producer"
"github.com/a8m/kinesis-producer/loggers/kplogrus"
"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/service/kinesis"
"github.com/openflagr/flagr/pkg/config"
"github.com/openflagr/flagr/swagger_gen/models"
"github.com/sirupsen/logrus"
)
var (
newKinesisProducer = producer.New
)
type kinesisRecorder struct {
producer *producer.Producer
options DataRecordFrameOptions
}
// NewKinesisRecorder creates a new Kinesis recorder
var NewKinesisRecorder = func() DataRecorder {
se, err := session.NewSession(aws.NewConfig())
if err != nil {
logrus.WithField("kinesis_error", err).Fatal("error creating aws session")
}
client := kinesis.New(se)
p := newKinesisProducer(&producer.Config{
StreamName: config.Config.RecorderKinesisStreamName,
Client: client,
BacklogCount: config.Config.RecorderKinesisBacklogCount,
MaxConnections: config.Config.RecorderKinesisMaxConnections,
FlushInterval: config.Config.RecorderKinesisFlushInterval,
BatchSize: config.Config.RecorderKinesisBatchSize,
BatchCount: config.Config.RecorderKinesisBatchCount,
AggregateBatchCount: config.Config.RecorderKinesisAggregateBatchCount,
AggregateBatchSize: config.Config.RecorderKinesisAggregateBatchSize,
Verbose: config.Config.RecorderKinesisVerbose,
Logger: &kplogrus.Logger{Logger: logrus.StandardLogger()},
})
p.Start()
go func() {
for err := range p.NotifyFailures() {
logrus.WithField("kinesis_error", err).Error("error pushing to kinesis")
}
}()
return &kinesisRecorder{
producer: p,
options: DataRecordFrameOptions{
Encrypted: false, // not implemented yet
FrameOutputMode: config.Config.RecorderFrameOutputMode,
},
}
}
func (k *kinesisRecorder) NewDataRecordFrame(r models.EvalResult) DataRecordFrame {
return DataRecordFrame{
evalResult: r,
options: k.options,
}
}
func (k *kinesisRecorder) AsyncRecord(r models.EvalResult) {
frame := k.NewDataRecordFrame(r)
output, err := frame.Output()
if err != nil {
logrus.WithField("err", err).Error("failed to generate data record frame for kinesis recorder")
return
}
err = k.producer.Put(output, frame.GetPartitionKey())
if err != nil {
logrus.WithField("kinesis_error", err).Error("error pushing to kinesis")
}
}