// client2server - Redpanda consumer // // Reads router-events from Redpanda and writes to SQLite for historical // queries. Started in background. package main import ( "context" "encoding/json" "log" "time" "github.com/twmb/franz-go/pkg/kgo" ) func startConsumer(ctx context.Context) { if kcl == nil { return } cl, err := kgo.NewClient( kgo.SeedBrokers(cfg.RedpandaBrokers...), kgo.ConsumerGroup("client2server-persistor"), kgo.ConsumeTopics("router-events"), kgo.DisableAutoCommit(), ) if err != nil { log.Printf("consumer: %v", err) return } defer cl.Close() log.Println("consumer: started, reading router-events") for { if ctx.Err() != nil { return } fetches := cl.PollFetches(ctx) if errs := fetches.Errors(); len(errs) > 0 { for _, e := range errs { log.Printf("consumer poll: %v", e.Err) } } fetches.EachRecord(func(rec *kgo.Record) { var ev RouterEvent if err := json.Unmarshal(rec.Value, &ev); err != nil { log.Printf("consumer decode: %v", err) return } // Ensure timestamps if ev.ReceivedAt.IsZero() { ev.ReceivedAt = rec.Timestamp } if err := saveEvent(ev); err != nil { log.Printf("consumer save: %v", err) } }) // Best-effort commit if err := cl.CommitUncommittedOffsets(ctx); err != nil { log.Printf("consumer commit: %v", err) } _ = time.Second } }