| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465 |
- // 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
- }
- }
|