Skip to content

Commit

Permalink
Merge pull request #29 from deepfence/switch-kafka-lib
Browse files Browse the repository at this point in the history
kafka: Replace confluentinc/confluent-kafka-go with segmentio/kafka-go
  • Loading branch information
vadorovsky authored Jul 10, 2022
2 parents cbc3c1b + b9dc1f4 commit 191e826
Show file tree
Hide file tree
Showing 5 changed files with 163 additions and 75 deletions.
5 changes: 3 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,12 @@ require (
github.com/aws/aws-sdk-go-v2 v1.16.4
github.com/aws/aws-sdk-go-v2/config v1.15.3
github.com/aws/aws-sdk-go-v2/service/s3 v1.26.3
github.com/confluentinc/confluent-kafka-go v1.8.2
github.com/foxcpp/go-mockdns v1.0.0
github.com/google/gopacket v1.1.19
github.com/google/uuid v1.3.0
github.com/inhies/go-bytesize v0.0.0-20210819104631-275770b98743
github.com/klauspost/compress v1.14.1
github.com/klauspost/compress v1.14.2
github.com/segmentio/kafka-go v0.4.32
github.com/spf13/cobra v1.4.0
gopkg.in/yaml.v3 v3.0.1
)
Expand All @@ -33,6 +33,7 @@ require (
github.com/inconshreveable/mousetrap v1.0.0 // indirect
github.com/kr/pretty v0.1.0 // indirect
github.com/miekg/dns v1.1.25 // indirect
github.com/pierrec/lz4/v4 v4.1.14 // indirect
github.com/spf13/pflag v1.0.5 // indirect
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550 // indirect
golang.org/x/net v0.0.0-20211209124913-491a49abca63 // indirect
Expand Down
26 changes: 19 additions & 7 deletions go.sum
Original file line number Diff line number Diff line change
@@ -1,4 +1,3 @@
github.com/aws/aws-sdk-go-v2 v1.16.2 h1:fqlCk6Iy3bnCumtrLz9r3mJ/2gUT0pJ0wLFVIdWh+JA=
github.com/aws/aws-sdk-go-v2 v1.16.2/go.mod h1:ytwTPBG6fXTZLxxeeCCWj2/EMYp/xDUgX+OET6TLNNU=
github.com/aws/aws-sdk-go-v2 v1.16.4 h1:swQTEQUyJF/UkEA94/Ga55miiKFoXmm/Zd67XHgmjSg=
github.com/aws/aws-sdk-go-v2 v1.16.4/go.mod h1:ytwTPBG6fXTZLxxeeCCWj2/EMYp/xDUgX+OET6TLNNU=
Expand Down Expand Up @@ -32,10 +31,10 @@ github.com/aws/aws-sdk-go-v2/service/sts v1.16.3 h1:cJGRyzCSVwZC7zZZ1xbx9m32UnrK
github.com/aws/aws-sdk-go-v2/service/sts v1.16.3/go.mod h1:bfBj0iVmsUyUg4weDB4NxktD9rDGeKSVWnjTnwbx9b8=
github.com/aws/smithy-go v1.11.2 h1:eG/N+CcUMAvsdffgMvjMKwfyDzIkjM6pfxMJ8Mzc6mE=
github.com/aws/smithy-go v1.11.2/go.mod h1:3xHYmszWVx2c0kIwQeEVf9uSm4fYZt67FBJnwub1bgM=
github.com/confluentinc/confluent-kafka-go v1.8.2 h1:PBdbvYpyOdFLehj8j+9ba7FL4c4Moxn79gy9cYKxG5E=
github.com/confluentinc/confluent-kafka-go v1.8.2/go.mod h1:u2zNLny2xq+5rWeTQjFHbDzzNuba4P1vo31r9r4uAdg=
github.com/cpuguy83/go-md2man/v2 v2.0.1/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/foxcpp/go-mockdns v1.0.0 h1:7jBqxd3WDWwi/6WhDvacvH1XsN3rOLXyHM1uhvIx6FI=
github.com/foxcpp/go-mockdns v1.0.0/go.mod h1:lgRN6+KxQBawyIghpnl5CezHFGS9VLzvtVlwxvzXTQ4=
github.com/google/go-cmp v0.5.7 h1:81/ik6ipDQS2aGcBfIN5dHDB36BwrStyeAQquSYCV4o=
Expand All @@ -50,23 +49,35 @@ github.com/inhies/go-bytesize v0.0.0-20210819104631-275770b98743 h1:X3Xxno5Ji8id
github.com/inhies/go-bytesize v0.0.0-20210819104631-275770b98743/go.mod h1:KrtyD5PFj++GKkFS/7/RRrfnRhAMGQwy75GLCHWrCNs=
github.com/jmespath/go-jmespath v0.4.0/go.mod h1:T8mJZnbsbmF+m6zOOFylbeCJqk5+pHWvzYPziyZiYoo=
github.com/jmespath/go-jmespath/internal/testify v1.5.1/go.mod h1:L3OGu8Wl2/fWfCI6z80xFu9LTZmf1ZRjMHUOPmWr69U=
github.com/klauspost/compress v1.14.1 h1:hLQYb23E8/fO+1u53d02A97a8UnsddcvYzq4ERRU4ds=
github.com/klauspost/compress v1.14.1/go.mod h1:/3/Vjq9QcHkK5uEr5lBEmyoZ1iFhe47etQ6QUkpK6sk=
github.com/klauspost/compress v1.14.2 h1:S0OHlFk/Gbon/yauFJ4FfJJF5V0fc5HbBTJazi28pRw=
github.com/klauspost/compress v1.14.2/go.mod h1:/3/Vjq9QcHkK5uEr5lBEmyoZ1iFhe47etQ6QUkpK6sk=
github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI=
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/miekg/dns v1.1.25 h1:dFwPR6SfLtrSwgDcIq2bcU/gVutB4sNApq2HBdqcakg=
github.com/miekg/dns v1.1.25/go.mod h1:bPDLeHnStXmXAq1m/Ch/hvfNHr14JKNPMBo3VZKjuso=
github.com/pierrec/lz4/v4 v4.1.14 h1:+fL8AQEZtz/ijeNnpduH0bROTu0O3NZAlPjQxGn8LwE=
github.com/pierrec/lz4/v4 v4.1.14/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
github.com/segmentio/kafka-go v0.4.32 h1:Ohr+9E+kDv/Ld2UPJN9hnKZRd2qgiqCmI8v2e1qlfLM=
github.com/segmentio/kafka-go v0.4.32/go.mod h1:JAPPIiY3MQIwVHj64CWOP0LsFFfQ7H0w69kuoxnMIS0=
github.com/spf13/cobra v1.4.0 h1:y+wJpx64xcgO1V+RcnwW0LEHxTKRi2ZDPSBjWnrg88Q=
github.com/spf13/cobra v1.4.0/go.mod h1:Wo4iy3BUC+X2Fybo0PDqwJIv3dNRiZLHQymsfxlB84g=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
github.com/spf13/pflag v1.0.5/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.7.1 h1:5TQK59W5E3v0r2duFAb7P95B6hEeOyEnHRa8MjYSMTY=
github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/xdg/scram v0.0.0-20180814205039-7eeb5667e42c h1:u40Z8hqBAAQyv+vATcGgV0YCnDjqSL7/q/JyPhhJSPk=
github.com/xdg/scram v0.0.0-20180814205039-7eeb5667e42c/go.mod h1:lB8K/P019DLNhemzwFU4jHLhdvlE6uDZjXFejJXr49I=
github.com/xdg/stringprep v1.0.0 h1:d9X0esnoa3dFsV0FG35rAT0RIhYFlPq7MiP+DW89La0=
github.com/xdg/stringprep v1.0.0/go.mod h1:Jhud4/sHMO4oL310DaZAKk9ZaJ08SJfe+sJh0HrGL1Y=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20190506204251-e1dfcc566284/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/crypto v0.0.0-20190923035154-9ee001bba392/go.mod h1:/lpIB1dKB+9EgE3H3cr1v9wB50oz8l4C4h62xy7jSTY=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550 h1:ObdrDkeb4kJdCP557AjRjq69pTHfNouLtWZG7j9rPN8=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
Expand All @@ -90,6 +101,7 @@ golang.org/x/sys v0.0.0-20211210111614-af8b64212486/go.mod h1:oPkhp1MJrh7nUepCBc
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk=
golang.org/x/text v0.3.6 h1:aRYxNxv6iGQlyVaZmk6ZgYEDa+Jg18DxebPSrd6bg1M=
golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20190907020128-2ca718005c18/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
Expand All @@ -103,7 +115,7 @@ gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15 h1:YR8cESwS4TdDjEe65xsg0ogR
gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ=
gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b h1:h8qDotaEPuJATrMmW04NCwg7v22aHH28wwpauUhK9Oo=
gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.0-20220512140231-539c8e751b99/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
23 changes: 14 additions & 9 deletions pkg/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,13 @@ package config

import (
"fmt"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/inhies/go-bytesize"
"io/ioutil"
"strings"
"time"

"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/inhies/go-bytesize"

"github.com/klauspost/compress/s2"
"gopkg.in/yaml.v3"
)
Expand Down Expand Up @@ -53,12 +55,13 @@ type S3PluginConfig struct {
}

type KafkaPluginConfig struct {
Brokers string
Brokers []string
ClientId string `yaml:"clientId,omitempty"`
Topic string `yaml:"topic,omitempty"`
MessageSize *bytesize.ByteSize `yaml:"messageSize,omitempty"`
Acks string `yaml:"acks,omitempty"`
FileSize *bytesize.ByteSize `yaml:"fileSize,omitempty"`
Timeout time.Duration `yaml:"timeout,omitempty"`
}

type PluginsConfig struct {
Expand All @@ -83,11 +86,12 @@ type S3OutputRawConfig struct {

type KafkaOutputRawConfig struct {
Brokers string
ClientId *string `yaml:"clientId,omitempty"`
Topic *string `yaml:"topic,omitempty"`
MessageSize *string `yaml:"messageSize,omitempty"`
Acks *string `yaml:"acks,omitempty"`
FileSize *string `yaml:"fileSize,omitempty"`
ClientId *string `yaml:"clientId,omitempty"`
Topic *string `yaml:"topic,omitempty"`
MessageSize *string `yaml:"messageSize,omitempty"`
Acks *string `yaml:"acks,omitempty"`
FileSize *string `yaml:"fileSize,omitempty"`
Timeout time.Duration `yaml:"timeout,omitempty"`
}

type PluginsRawConfig struct {
Expand Down Expand Up @@ -289,12 +293,13 @@ func populateKafkaConfig(rawConfig RawConfig) (*KafkaPluginConfig, error) {
}

return &KafkaPluginConfig{
Brokers: rawConfig.Output.Plugins.Kafka.Brokers,
Brokers: strings.Split(rawConfig.Output.Plugins.Kafka.Brokers, ","),
ClientId: clientId,
Topic: topic,
MessageSize: messageSize,
Acks: acks,
FileSize: fileSize,
Timeout: rawConfig.Output.Plugins.Kafka.Timeout,
}, nil
}

Expand Down
66 changes: 38 additions & 28 deletions pkg/plugins/kafka/kafka.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,18 @@ package kafka

import (
"context"
"fmt"
"github.com/confluentinc/confluent-kafka-go/kafka"
"log"
"time"

"github.com/deepfence/PacketStreamer/pkg/config"
"github.com/deepfence/PacketStreamer/pkg/file"
"github.com/google/uuid"
"log"
kafka "github.com/segmentio/kafka-go"
)

type KafkaProducer interface {
Produce(msg *kafka.Message, deliveryChan chan kafka.Event) error
Close()
type KafkaWriter interface {
WriteMessages(ctx context.Context, msgs ...kafka.Message) error
Close() error
}

type File struct {
Expand All @@ -25,8 +26,19 @@ func (f *File) newBuffer(size int) {
f.Buffer = make([]byte, 0, size)
}

type IdGenerator interface {
Generate() string
}

type FileIdGenerator struct{}

func (g *FileIdGenerator) Generate() string {
return uuid.New().String()
}

type Plugin struct {
Producer KafkaProducer
Writer KafkaWriter
IdGenerator IdGenerator
Topic string
MessageSize int
CloseChan chan bool
Expand All @@ -35,18 +47,20 @@ type Plugin struct {
}

func NewPlugin(config *config.KafkaPluginConfig) (*Plugin, error) {
producer, err := kafka.NewProducer(&kafka.ConfigMap{
"bootstrap.servers": config.Brokers,
"client.id": config.ClientId,
"acks": config.Acks,
})

if err != nil {
return nil, fmt.Errorf("error creating kafka producer, %w", err)
dialer := &kafka.Dialer{
Timeout: 10 * time.Second,
DualStack: true,
TLS: nil,
}
writer := kafka.NewWriter(kafka.WriterConfig{
Brokers: config.Brokers,
Topic: config.Topic,
Balancer: &kafka.Hash{},
Dialer: dialer,
})

return &Plugin{
Producer: producer,
Writer: writer,
Topic: config.Topic,
MessageSize: int(*config.MessageSize),
FileSize: uint64(*config.FileSize),
Expand All @@ -67,8 +81,8 @@ func (p *Plugin) newFile(id string, messageSize int) {
func (p *Plugin) Start(ctx context.Context) chan<- string {
inputChan := make(chan string)
go func() {
defer p.Producer.Close()
p.newFile(generateFileId(), p.MessageSize)
defer p.Writer.Close()
p.newFile(p.IdGenerator.Generate(), p.MessageSize)

for {
select {
Expand Down Expand Up @@ -103,7 +117,7 @@ func (p *Plugin) Start(ctx context.Context) chan<- string {
}

if p.CurrentFile.Sent >= p.FileSize {
p.newFile(generateFileId(), p.MessageSize)
p.newFile(p.IdGenerator.Generate(), p.MessageSize)
} else {
p.CurrentFile.newBuffer(p.MessageSize)
}
Expand Down Expand Up @@ -132,17 +146,13 @@ func (p *Plugin) cleanup() {
}

func (p *Plugin) flush() error {
err := p.Producer.Produce(&kafka.Message{
TopicPartition: kafka.TopicPartition{Topic: &p.Topic, Partition: kafka.PartitionAny},
Value: p.CurrentFile.Buffer,
Key: []byte(p.CurrentFile.Id),
}, nil)
err := p.Writer.WriteMessages(context.Background(), kafka.Message{
Topic: p.Topic,
Key: []byte(p.CurrentFile.Id),
Value: p.CurrentFile.Buffer,
})

p.CurrentFile.Sent += uint64(len(p.CurrentFile.Buffer))

return err
}

func generateFileId() string {
return uuid.New().String()
}
Loading

0 comments on commit 191e826

Please sign in to comment.