Skip to content

Примеры 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> -L

Producer

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 и дедупликацию ретраев в рамках сессии.