Примеры Kafka-клиентов
Адрес брокера — SERVER:9092. Если включена аутентификация ([[auth.users]] в serve.toml) — добавьте SASL/PLAIN, как в примерах ниже; без неё параметры sasl_*/security.protocol уберите. Для доступа с других машин у брокера должен быть настроен kafka_advertised_addr (см. конфигурацию).
Проверка соединения
bash
kcat -b SERVER:9092 -L # metadata: брокер + топики
echo "hello" | kcat -b SERVER:9092 -t demo -P # produce
kcat -b SERVER:9092 -t demo -C -o beginning -e # consume с началаbash
kcat -b SERVER:9092 \
-X security.protocol=SASL_PLAINTEXT -X sasl.mechanism=PLAIN \
-X sasl.username=app -X sasl.password=<pass> -LProducer
python
# pip install confluent-kafka (librdkafka — проверено e2e с Roxa)
from confluent_kafka import Producer
p = Producer({
"bootstrap.servers": "SERVER:9092",
# при включённой аутентификации:
"security.protocol": "SASL_PLAINTEXT",
"sasl.mechanism": "PLAIN",
"sasl.username": "app",
"sasl.password": "<pass>",
})
p.produce("demo", key="order-1", value=b"hello roxa")
p.flush() # ack придёт после durable-записи пачки в хранилищеgo
// go get github.com/twmb/franz-go/pkg/kgo
package main
import (
"context"
"github.com/twmb/franz-go/pkg/kgo"
"github.com/twmb/franz-go/pkg/sasl/plain"
)
func main() {
cl, err := kgo.NewClient(
kgo.SeedBrokers("SERVER:9092"),
// при включённой аутентификации:
kgo.SASL(plain.Auth{User: "app", Pass: "<pass>"}.AsMechanism()),
)
if err != nil { panic(err) }
defer cl.Close()
rec := &kgo.Record{Topic: "demo", Key: []byte("order-1"), Value: []byte("hello roxa")}
if err := cl.ProduceSync(context.Background(), rec).FirstErr(); err != nil {
panic(err)
}
}java
// Java-клиент говорит на той же спецификации протокола, но в e2e-тестах Roxa
// пока не гоняется — проверьте свой сценарий на демо-стенде.
Properties props = new Properties();
props.put("bootstrap.servers", "SERVER:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// при включённой аутентификации:
props.put("security.protocol", "SASL_PLAINTEXT");
props.put("sasl.mechanism", "PLAIN");
props.put("sasl.jaas.config",
"org.apache.kafka.common.security.plain.PlainLoginModule required " +
"username=\"app\" password=\"<pass>\";");
try (var producer = new KafkaProducer<String, String>(props)) {
producer.send(new ProducerRecord<>("demo", "order-1", "hello roxa")).get();
}Consumer (группа с ребалансом)
python
from confluent_kafka import Consumer
c = Consumer({
"bootstrap.servers": "SERVER:9092",
"group.id": "billing",
"auto.offset.reset": "earliest",
# + sasl-параметры как у продюсера, если включена аутентификация
})
c.subscribe(["demo"])
while True:
msg = c.poll(1.0)
if msg is None:
continue
if msg.error():
raise Exception(msg.error())
print(msg.key(), msg.value())
c.commit(msg)go
cl, _ := kgo.NewClient(
kgo.SeedBrokers("SERVER:9092"),
kgo.ConsumerGroup("billing"),
kgo.ConsumeTopics("demo"),
)
defer cl.Close()
for {
fetches := cl.PollFetches(context.Background())
fetches.EachRecord(func(r *kgo.Record) {
fmt.Printf("%s = %s\n", r.Key, r.Value)
})
cl.CommitUncommittedOffsets(context.Background())
}Идемпотентный продюсер
Включается штатно (enable.idempotence=true в librdkafka-клиентах) — Roxa поддерживает InitProducerId и дедупликацию ретраев в рамках сессии.