/
writer.go
61 lines (52 loc) · 1.42 KB
/
writer.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
package meta
import (
"context"
"fmt"
"github.com/segmentio/kafka-go"
)
// KafkaWriter is an abstraction over kafka-go client implementation
type KafkaWriter interface {
WriteMessages(context.Context, ...kafka.Message) error
Close() error
Stats() kafka.WriterStats
}
// Writer will be used to write send data to kafka topic
type Writer struct {
client KafkaWriter
bufferSize int
bufferedMessages []kafka.Message
}
// NewWriter returns a instance for writer used over kafka client
func NewWriter(w KafkaWriter, buffSize int) *Writer {
return &Writer{
client: w,
bufferSize: buffSize,
bufferedMessages: make([]kafka.Message, 0),
}
}
// Write push messages to kafka
// this will throw an error if connection was closed in the middle of write
func (w *Writer) Write(protobufkey []byte, protobuf []byte) error {
msg := kafka.Message{
Key: protobufkey,
Value: protobuf,
}
w.bufferedMessages = append(w.bufferedMessages, msg)
var err error
if len(w.bufferedMessages) >= w.bufferSize {
err = w.Flush()
}
return err
}
// Flush will push all the queued up messages to kafka
func (w *Writer) Flush() error {
var err error
if len(w.bufferedMessages) > 0 {
err = w.client.WriteMessages(context.Background(), w.bufferedMessages...)
if err == nil {
w.bufferedMessages = make([]kafka.Message, 0)
fmt.Println("Published metadata for", len(w.bufferedMessages), "specs")
}
}
return err
}