Skip to content

Commit 950b088

Browse files
committed
create kafka exporter addon
1 parent 56827cb commit 950b088

7 files changed

Lines changed: 620 additions & 2 deletions

File tree

addons/confluent-kafka/go.mod

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
module github.com/spdeepak/flowtracker/confluent-kafka
2+
3+
go 1.24
4+
5+
require (
6+
github.com/confluentinc/confluent-kafka-go v1.9.2
7+
github.com/spdeepak/flowtracker v0.0.3
8+
)

addons/confluent-kafka/go.sum

Lines changed: 213 additions & 0 deletions
Large diffs are not rendered by default.

addons/confluent-kafka/kafka.go

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,114 @@
1+
package confluent_kafka
2+
3+
import (
4+
"encoding/json"
5+
"fmt"
6+
"log"
7+
8+
"github.com/confluentinc/confluent-kafka-go/kafka"
9+
"github.com/spdeepak/flowtracker"
10+
)
11+
12+
// Config holds the setup parameters.
13+
// You must provide EITHER Producer OR KafkaConfigMap.
14+
type Config struct {
15+
// Topic is the destination Kafka topic name (Required).
16+
Topic string
17+
18+
// Producer allows you to pass an existing Kafka Producer.
19+
// If this is set, KafkaConfigMap is ignored.
20+
Producer *kafka.Producer
21+
22+
// KafkaConfigMap allows you to configure a new Producer.
23+
// Use this if you don't have an existing client.
24+
// Example: &kafka.ConfigMap{"bootstrap.servers": "localhost:9092"}
25+
KafkaConfigMap *kafka.ConfigMap
26+
}
27+
28+
// KafkaExporter implements the flowtracker.Exporter interface.
29+
type KafkaExporter struct {
30+
producer *kafka.Producer
31+
topic string
32+
// isOwned tracks if this exporter created the producer (and thus should close it).
33+
isOwned bool
34+
}
35+
36+
// New creates a new KafkaExporter.
37+
func New(cfg Config) (*KafkaExporter, error) {
38+
if cfg.Topic == "" {
39+
return nil, fmt.Errorf("kafka exporter: topic is required")
40+
}
41+
42+
var p *kafka.Producer
43+
var err error
44+
isOwned := false
45+
46+
if cfg.Producer != nil {
47+
// Use the user-provided producer
48+
p = cfg.Producer
49+
} else if cfg.KafkaConfigMap != nil {
50+
// Initialize a new producer
51+
p, err = kafka.NewProducer(cfg.KafkaConfigMap)
52+
if err != nil {
53+
return nil, fmt.Errorf("kafka exporter: failed to create producer: %w", err)
54+
}
55+
isOwned = true
56+
57+
// Background goroutine to handle delivery reports.
58+
// Essential for confluent-kafka-go to prevent local queue filling up.
59+
go func() {
60+
for e := range p.Events() {
61+
switch ev := e.(type) {
62+
case *kafka.Message:
63+
if ev.TopicPartition.Error != nil {
64+
log.Printf("FlowTracker Kafka Error: %v\n", ev.TopicPartition.Error)
65+
}
66+
}
67+
}
68+
}()
69+
} else {
70+
return nil, fmt.Errorf("kafka exporter: must provide either Producer or KafkaConfigMap")
71+
}
72+
73+
return &KafkaExporter{
74+
producer: p,
75+
topic: cfg.Topic,
76+
isOwned: isOwned,
77+
}, nil
78+
}
79+
80+
// Export sends the trace to Kafka.
81+
func (k *KafkaExporter) Export(tr *flowtracker.Trace) {
82+
// Serialize Trace to JSON
83+
payload, err := json.Marshal(tr)
84+
if err != nil {
85+
log.Printf("FlowTracker: failed to marshal trace: %v", err)
86+
return
87+
}
88+
89+
// Construct the Kafka Message
90+
msg := &kafka.Message{
91+
TopicPartition: kafka.TopicPartition{Topic: &k.topic, Partition: kafka.PartitionAny},
92+
Value: payload,
93+
// We use the TraceID as the Key. This ensures that if you update this logic
94+
// to stream updates, all spans for the same trace go to the same partition.
95+
Key: []byte(tr.TraceID),
96+
}
97+
98+
// Produce is asynchronous. We rely on the background event loop (started in New) to handle errors.
99+
err = k.producer.Produce(msg, nil)
100+
if err != nil {
101+
log.Printf("FlowTracker: failed to produce message: %v", err)
102+
}
103+
}
104+
105+
// Close flushes the producer if we own it.
106+
// Note: The main flowtracker library doesn't call Close(), but you can call this
107+
// manually in your main.go shutdown hook.
108+
func (k *KafkaExporter) Close() {
109+
if k.isOwned && k.producer != nil {
110+
// Wait up to 5 seconds for outstanding messages to be delivered
111+
k.producer.Flush(5000)
112+
k.producer.Close()
113+
}
114+
}
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
version: "3"
2+
services:
3+
zookeeper:
4+
image: "zookeeper:latest"
5+
environment:
6+
- ZOO_CFG_EXTRA=autopurge.snapRetainCount=3
7+
- ZOO_CFG_EXTRA=autopurge.purgeInterval=24
8+
ports:
9+
- "2181:2181"
10+
restart: unless-stopped
11+
12+
kafka:
13+
image: wurstmeister/kafka
14+
ports:
15+
- "9092:9092"
16+
- "19092:19092"
17+
volumes:
18+
- ./test.json:/test.json
19+
environment:
20+
KAFKA_ADVERTISED_HOST_NAME: localhost
21+
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
22+
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:19092
23+
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_HOST://localhost:19092
24+
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
25+
KAFKA_CREATE_TOPICS: "flowtracker.logs:1:1"
26+
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
27+
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
28+
depends_on:
29+
- zookeeper
30+
restart: unless-stopped

examples/confluent-kafka/kafka.go

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
package main
2+
3+
import (
4+
"log/slog"
5+
"net/http"
6+
"os"
7+
8+
"github.com/confluentinc/confluent-kafka-go/kafka"
9+
"github.com/spdeepak/flowtracker"
10+
confluentkafka "github.com/spdeepak/flowtracker/confluent-kafka"
11+
"github.com/spdeepak/flowtracker/examples"
12+
)
13+
14+
func main() {
15+
mux := http.NewServeMux()
16+
mux.HandleFunc("/", examples.Handler)
17+
18+
config := confluentkafka.Config{
19+
Topic: "flowtracker.logs",
20+
KafkaConfigMap: &kafka.ConfigMap{
21+
"bootstrap.servers": "localhost:9092",
22+
"client.id": "flowtracker-client",
23+
"acks": "all",
24+
},
25+
}
26+
exporter, err := confluentkafka.New(config)
27+
if err != nil {
28+
slog.Error("Error during kafka exporter creation", slog.Any("error", err))
29+
os.Exit(1)
30+
}
31+
mw := flowtracker.NewMiddleware(flowtracker.WithExporter(exporter))
32+
33+
http.ListenAndServe(":8080", mw(mux))
34+
}

examples/go.mod

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,14 @@ module github.com/spdeepak/flowtracker/examples
22

33
go 1.24
44

5-
require github.com/spdeepak/flowtracker v0.0.3
5+
require (
6+
github.com/spdeepak/flowtracker v0.0.3
7+
github.com/spdeepak/flowtracker/confluent-kafka v0.0.1
8+
)
69

7-
replace github.com/spdeepak/flowtracker => ../
10+
require github.com/confluentinc/confluent-kafka-go v1.9.2 // indirect
11+
12+
replace (
13+
github.com/spdeepak/flowtracker => ../
14+
github.com/spdeepak/flowtracker/confluent-kafka => ../addons/confluent-kafka/
15+
)

0 commit comments

Comments
 (0)