Table of contents
Part of nurago, a collection of independent Go packages for backend services.
import "github.com/tecnickcom/nurago/pkg/kafka"
Package kafka provides a pure-Go API for producing and consuming Apache Kafka messages. It requires no CGO and no system librdkafka installation.
Built on github.com/segmentio/kafka-go, it exposes a producer/consumer interface:
NewProducer+Producer.Send/Producer.SendDataNewConsumer+Consumer.Receive/Consumer.ReceiveDataNewConsumer+Consumer.FetchMessage/Consumer.CommitMessages
Functional options tune the session timeout, start offset, required acks, batching, and custom codecs.
Delivery Semantics
Producer writes wait for broker acknowledgment from the full in-sync replica
set by default (kafka.RequireAll); tune this with WithRequiredAcks.
On the consumer side, when a consumer group is configured:
Consumer.ReceiveandConsumer.ReceiveDataare at-most-once: the offset is committed as soon as the message is read, before the caller processes it, so a crash or decode failure after the read permanently skips the message.Consumer.FetchMessage+Consumer.CommitMessagesare at-least-once: the offset is committed only when the caller explicitly acknowledges the message after successful processing.
When no consumer group is configured (empty groupID), offsets are never
committed: Consumer.Receive and Consumer.FetchMessage behave identically,
reading always starts from the earliest available offset, and
Consumer.CommitMessages returns an error.
Message Encoding and Decoding
Typed payload methods use configurable codec hooks:
DefaultMessageEncodeFuncpowersProducer.SendDataDefaultMessageDecodeFuncpowersConsumer.ReceiveData
Both defaults use github.com/tecnickcom/nurago/pkg/encode. Replace them via
WithMessageEncodeFunc and WithMessageDecodeFunc to add custom wire formats,
encryption, compression, or schema validation.
Errors
Configuration problems are reported at construction time with errors matching
the exported sentinels ErrInvalidOptions, ErrNilEncodeFunc, and
ErrNilDecodeFunc. After Consumer.Close, the receive methods return errors
matching ErrConsumerClosed. Match the sentinels with errors.Is.
When To Use
- You want a Kafka client that cross-compiles and needs no system library.
- Message payloads are structs and should be encoded and decoded consistently.
- At-least-once delivery requires fetching and committing offsets separately.
Example
queue := &exampleQueue{}
producer, err := kafka.NewProducer(
[]string{"127.0.0.1:9092"},
"events",
kafka.WithKafkaWriter(queue),
)
if err != nil {
fmt.Println(err)
return
}
defer func() { _ = producer.Close() }()
ctx := context.TODO()
err = producer.Send(ctx, []byte("raw payload"))
if err != nil {
fmt.Println(err)
return
}
consumer, err := kafka.NewConsumer(
[]string{"127.0.0.1:9092"},
"events",
"example-group",
kafka.WithKafkaReader(queue),
)
if err != nil {
fmt.Println(err)
return
}
defer func() { _ = consumer.Close() }()
msg, err := consumer.Receive(ctx)
fmt.Println(string(msg), err)
// Output:
// raw payload <nil>
Full source is in example_kafka_test.go. More runnable examples are on pkg.go.dev.
Dependencies
Importing this package pulls 3 external modules:
github.com/klauspost/compressgithub.com/pierrec/lz4/v4github.com/segmentio/kafka-go