11package main
22
33import (
4+ "bytes"
45 "context"
56 "log"
67 "net"
78 "os"
89 "os/signal"
910 "runtime"
11+ "strings"
1012 "syscall"
13+ "time"
14+
15+ "github.com/aws/aws-sdk-go-v2/config"
16+ "github.com/netsampler/goflow2/v2/decoders/sflow"
17+ "github.com/twmb/franz-go/pkg/kgo"
18+ "github.com/twmb/franz-go/pkg/kversion"
19+ "github.com/twmb/franz-go/pkg/sasl/aws"
1120)
1221
1322type packet struct {
@@ -44,9 +53,49 @@ func main() {
4453
4554 packets := make (chan packet , 1024 )
4655
56+ kafkaBrokers := os .Getenv ("KAFKA_BROKERS" )
57+ if kafkaBrokers == "" {
58+ kafkaBrokers = "127.0.0.1:19092"
59+ }
60+
61+ seeds := kgo .SeedBrokers (strings .Split (kafkaBrokers , "," )... )
62+ opts := []kgo.Opt {
63+ seeds ,
64+ kgo .AllowAutoTopicCreation (),
65+ kgo .RequiredAcks (kgo .AllISRAcks ()),
66+ kgo .ProducerBatchCompression (kgo .SnappyCompression ()),
67+ kgo .ProducerLinger (1 * time .Second ),
68+ kgo .MaxVersions (kversion .V2_8_0 ()),
69+ }
70+
71+ if strings .ToLower (os .Getenv ("KAFKA_AUTH_IAM_ENABLED" )) == "true" {
72+ awsCfg , err := config .LoadDefaultConfig (ctx )
73+ if err != nil {
74+ log .Fatalf ("failed to load aws config: %v" , err )
75+ }
76+ opts = append (opts , kgo .SASL (aws .ManagedStreamingIAM (func (ctx context.Context ) (aws.Auth , error ) {
77+ creds , err := awsCfg .Credentials .Retrieve (ctx )
78+ if err != nil {
79+ return aws.Auth {}, err
80+ }
81+ return aws.Auth {
82+ AccessKey : creds .AccessKeyID ,
83+ SecretKey : creds .SecretAccessKey ,
84+ SessionToken : creds .SessionToken ,
85+ }, nil
86+ })))
87+ opts = append (opts , kgo .DialTLS ())
88+ }
89+
90+ kafkaClient , err := kgo .NewClient (opts ... )
91+ if err != nil {
92+ log .Fatalf ("kafka client: %v" , err )
93+ }
94+ defer kafkaClient .Close ()
95+
4796 workerCount := runtime .NumCPU ()
4897 for i := 0 ; i < workerCount ; i ++ {
49- go ingestWorker (ctx , i , packets )
98+ go ingestWorker (ctx , i , packets , kafkaClient )
5099 }
51100
52101 readLoop (ctx , conn , packets )
@@ -86,7 +135,7 @@ func readLoop(ctx context.Context, conn *net.UDPConn, out chan<- packet) {
86135 }
87136}
88137
89- func ingestWorker (ctx context.Context , id int , in <- chan packet ) {
138+ func ingestWorker (ctx context.Context , id int , in <- chan packet , kafkaClient * kgo. Client ) {
90139 log .Printf ("worker %d started" , id )
91140 for {
92141 select {
@@ -98,13 +147,40 @@ func ingestWorker(ctx context.Context, id int, in <-chan packet) {
98147 log .Printf ("worker %d channel closed, exiting" , id )
99148 return
100149 }
101- ingestPacket (id , p )
150+ ingestPacket (ctx , id , p , kafkaClient )
102151 }
103152 }
104153}
105154
106- // replace this with "write to DB / queue / whatever"
107- func ingestPacket (workerID int , p packet ) {
108- log .Printf ("worker %d: %d bytes from %s" , workerID , len (p .data ), p .addr .String ())
109- // parse p.data and do your real ingestion here
155+ func ingestPacket (ctx context.Context , workerID int , p packet , kafkaClient * kgo.Client ) {
156+ // we need to check this a valid sflow packet before sending to kafka
157+ var msg sflow.Packet
158+ err := sflow .DecodeMessage (bytes .NewBuffer (p .data ), & msg )
159+ if err != nil {
160+ log .Printf ("worker %d: sflow decode error: %v" , workerID , err )
161+ return
162+ }
163+
164+ var hasFlowSample bool
165+ for _ , sample := range msg .Samples {
166+ if _ , ok := sample .(sflow.FlowSample ); ok {
167+ hasFlowSample = true
168+ break
169+ }
170+ }
171+
172+ if ! hasFlowSample {
173+ return // skip packets without flow samples
174+ }
175+
176+ rec := & kgo.Record {
177+ Topic : "flows_raw_devnet" ,
178+ Value : p .data ,
179+ }
180+
181+ kafkaClient .Produce (ctx , rec , func (r * kgo.Record , err error ) {
182+ if err != nil {
183+ log .Printf ("worker %d: kafka produce error: %v" , workerID , err )
184+ }
185+ })
110186}
0 commit comments