GitHub - samber/slog-kafka: 🚨 slog: Kafka handler (original) (raw)
slog: Kafka handler
A Kafka Handler for slog Go library.
See also:
- slog-multi:
slog.Handlerchaining, fanout, routing, failover, load balancing... - slog-formatter:
slogattribute formatting - slog-sampling:
slogsampling policy - slog-mock:
slog.Handlerfor test purposes
HTTP middlewares:
- slog-gin: Gin middleware for
sloglogger - slog-echo: Echo middleware for
sloglogger - slog-fiber: Fiber middleware for
sloglogger - slog-chi: Chi middleware for
sloglogger - slog-http:
net/httpmiddleware forsloglogger
Loggers:
- slog-zap: A
sloghandler forZap - slog-zerolog: A
sloghandler forZerolog - slog-logrus: A
sloghandler forLogrus
Log sinks:
- slog-datadog: A
sloghandler forDatadog - slog-betterstack: A
sloghandler forBetterstack - slog-rollbar: A
sloghandler forRollbar - slog-loki: A
sloghandler forLoki - slog-sentry: A
sloghandler forSentry - slog-syslog: A
sloghandler forSyslog - slog-logstash: A
sloghandler forLogstash - slog-fluentd: A
sloghandler forFluentd - slog-graylog: A
sloghandler forGraylog - slog-quickwit: A
sloghandler forQuickwit - slog-slack: A
sloghandler forSlack - slog-telegram: A
sloghandler forTelegram - slog-mattermost: A
sloghandler forMattermost - slog-microsoft-teams: A
sloghandler forMicrosoft Teams - slog-webhook: A
sloghandler forWebhook - slog-kafka: A
sloghandler forKafka - slog-nats: A
sloghandler forNATS - slog-parquet: A
sloghandler forParquet+Object Storage - slog-channel: A
sloghandler for Go channels
🚀 Install
go get github.com/samber/slog-kafka/v2
Compatibility: go >= 1.21
No breaking changes will be made to exported APIs before v3.0.0.
💡 Usage
GoDoc: https://pkg.go.dev/github.com/samber/slog-kafka/v2
Handler options
type Option struct { // log level (default: debug) Level slog.Leveler
// Kafka Writer
KafkaWriter *kafka.Writer
Timeout time.Duration // default: 60s
// optional: customize Kafka event builder
Converter Converter
// optional: custom marshaler
Marshaler func(v any) ([]byte, error)
// optional: fetch attributes from context
AttrFromContext []func(ctx context.Context) []slog.Attr
// optional: see slog.HandlerOptions
AddSource bool
ReplaceAttr func(groups []string, a slog.Attr) slog.Attr}
Other global parameters:
slogkafka.SourceKey = "source" slogkafka.ContextKey = "extra" slogkafka.RequestKey = "request" slogkafka.ErrorKeys = []string{"error", "err"} slogkafka.RequestIgnoreHeaders = false
Supported attributes
The following attributes are interpreted by slogkafka.DefaultConverter:
| Atribute name | slog.Kind | Underlying type |
|---|---|---|
| "user" | group (see below) | |
| "error" | any | error |
| "request" | any | *http.Request |
| other attributes | * |
Other attributes will be injected in extra field.
Users must be of type slog.Group. Eg:
slog.Group("user", slog.String("id", "user-123"), slog.String("username", "samber"), slog.Time("created_at", time.Now()), )
Example
import ( "context" "fmt" "time"
slogkafka "github.com/samber/slog-kafka/v2"
"github.com/segmentio/kafka-go"
"log/slog")
func main() { // docker-compose up -d
uri := "127.0.0.1:9092"
dialer := &kafka.Dialer{
Timeout: 10 * time.Second,
DualStack: true,
}
conn, err := dialer.DialContext(context.Background(), "tcp", uri)
if err != nil {
panic(err)
}
err = conn.CreateTopics(kafka.TopicConfig{
Topic: "logs",
NumPartitions: 12,
ReplicationFactor: 1,
})
if err != nil {
panic(err)
}
writer := kafka.NewWriter(kafka.WriterConfig{
Brokers: []string{uri},
Topic: "logs",
Dialer: dialer,
Async: true, // !
Balancer: &kafka.Hash{},
MaxAttempts: 3,
Logger: kafka.LoggerFunc(func(msg string, args ...interface{}) {
fmt.Printf(msg+"\n", args...)
}),
ErrorLogger: kafka.LoggerFunc(func(msg string, args ...interface{}) {
fmt.Printf(msg+"\n", args...)
}),
})
defer writer.Close()
defer conn.Close()
logger := slog.New(slogkafka.Option{Level: slog.LevelDebug, KafkaWriter: writer}.NewKafkaHandler())
logger = logger.With("release", "v1.0.0")
logger.
With(
slog.Group("user",
slog.String("id", "user-123"),
slog.Time("created_at", time.Now()),
),
).
With("error", fmt.Errorf("an error")).
Error("a message")}
Kafka message:
{ "level": "ERROR", "logger": "samber/slog-kafka", "message": "a message", "timestamp": "2023-04-30T01:33:21.676768Z", "error": { "error": "an error", "kind": "*errors.errorString", "stack": null }, "extra": { "release": "v1.0.0" }, "user": { "created_at": "2023-04-30T01:33:21.676704Z", "id": "user-123" } }
Tracing
Import the samber/slog-otel library.
import ( slogkafka "github.com/samber/slog-kafka" slogotel "github.com/samber/slog-otel" "go.opentelemetry.io/otel/sdk/trace" )
func main() { tp := trace.NewTracerProvider( trace.WithSampler(trace.AlwaysSample()), ) tracer := tp.Tracer("hello/world")
ctx, span := tracer.Start(context.Background(), "foo")
defer span.End()
span.AddEvent("bar")
logger := slog.New(
slogkafka.Option{
// ...
AttrFromContext: []func(ctx context.Context) []slog.Attr{
slogotel.ExtractOtelAttrFromContext([]string{"tracing"}, "trace_id", "span_id"),
},
}.NewKafkaHandler(),
)
logger.ErrorContext(ctx, "a message")}
🤝 Contributing
- Ping me on twitter @samuelberthe (DMs, mentions, whatever :))
- Fork the project
- Fix open issues or request new features
Don't hesitate ;)
Install some dev dependencies
make tools
Run tests
make test
or
make watch-test
👤 Contributors
💫 Show your support
Give a ⭐️ if this project helped you!
📝 License
Copyright © 2023 Samuel Berthe.
This project is MIT licensed.