consumer.go 1.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465
  1. // client2server - Redpanda consumer
  2. //
  3. // Reads router-events from Redpanda and writes to SQLite for historical
  4. // queries. Started in background.
  5. package main
  6. import (
  7. "context"
  8. "encoding/json"
  9. "log"
  10. "time"
  11. "github.com/twmb/franz-go/pkg/kgo"
  12. )
  13. func startConsumer(ctx context.Context) {
  14. if kcl == nil {
  15. return
  16. }
  17. cl, err := kgo.NewClient(
  18. kgo.SeedBrokers(cfg.RedpandaBrokers...),
  19. kgo.ConsumerGroup("client2server-persistor"),
  20. kgo.ConsumeTopics("router-events"),
  21. kgo.DisableAutoCommit(),
  22. )
  23. if err != nil {
  24. log.Printf("consumer: %v", err)
  25. return
  26. }
  27. defer cl.Close()
  28. log.Println("consumer: started, reading router-events")
  29. for {
  30. if ctx.Err() != nil {
  31. return
  32. }
  33. fetches := cl.PollFetches(ctx)
  34. if errs := fetches.Errors(); len(errs) > 0 {
  35. for _, e := range errs {
  36. log.Printf("consumer poll: %v", e.Err)
  37. }
  38. }
  39. fetches.EachRecord(func(rec *kgo.Record) {
  40. var ev RouterEvent
  41. if err := json.Unmarshal(rec.Value, &ev); err != nil {
  42. log.Printf("consumer decode: %v", err)
  43. return
  44. }
  45. // Ensure timestamps
  46. if ev.ReceivedAt.IsZero() {
  47. ev.ReceivedAt = rec.Timestamp
  48. }
  49. if err := saveEvent(ev); err != nil {
  50. log.Printf("consumer save: %v", err)
  51. }
  52. })
  53. // Best-effort commit
  54. if err := cl.CommitUncommittedOffsets(ctx); err != nil {
  55. log.Printf("consumer commit: %v", err)
  56. }
  57. _ = time.Second
  58. }
  59. }