Kafka-клиенты в разных языках: драйверы, родословные, различия

Клиент Kafka делает куда больше, чем «подключиться»: партиционирование, батчинг, сжатие, управление offset, ребаланс, идемпотентность, транзакции — и всё это разные библиотеки реализуют по-разному. Три родословные (референсный JVM-клиент, обёртки над librdkafka на C, чистые реализации протокола), почему у Go сразу несколько драйверов (franz-go/sarama/segmentio/confluent) и чем они отличаются, и что из этого следует для фич, производительности и деплоя. С живыми примерами produce/consume на Java, Go, Rust, Python, C#, C++, Scala

«Подключиться к Kafka» звучит как деталь, но клиентская библиотека — это половина системы: именно она решает, в какую партицию уйдёт сообщение, как оно забатчится и сожмётся, когда закоммитится offset, как поведёт себя consumer group при ребалансе, поддерживаются ли идемпотентный producer и транзакции. И тут выясняется, что «клиент Kafka» — не один клиент: у каждого языка их несколько, они разной родословной, с разным набором фич, производительностью и способом сборки. Выбор драйвера — реальное архитектурное решение, а не строчка в go.mod.

Это восьмая статья серии «Kafka: глубокое погружение». Опирается на модель партиционирования из первой статьи, механику ребаланса из статьи про consumer groups и разбор exactly-once из статьи про транзакции — здесь всё это собирается в одну картину: как одна и та же механика по-разному реализована в 11 драйверах на 7 языках, прогнанных живьём против общего кластера с одинаковым сценарием. Код всех 11 драйверов — в открытом виде, каталог kafka/clients репозитория digital-cookbook.

Клиенты Kafka в разных языках: три родословные партиционера дают три раскладки ключей

В статье

Что делает клиент

Со стороны кажется, что клиент Kafka — это сокет плюс сериализация. На деле клиентская библиотека закрывает собой большую часть протокольной сложности, которую брокер сознательно оставил на клиентской стороне:

  • Producer: выбирает партицию для записи (по ключу через партиционер либо round-robin/sticky, если ключа нет — подробный разбор партиционера на живых данных был в первой статье), батчирует записи (linger.ms/batch.size), сжимает батч целиком, ждёт нужный уровень acks, при необходимости включает идемпотентность (PID + sequence-номер на партицию) и ретраит без дублей.
  • Consumer: вступает в consumer group, договаривается с координатором о назначении партиций, переживает ребаланс — eager (stop-the-world) или cooperative-sticky (инкрементальный), — коммитит offset автоматически или вручную, соблюдает isolation level (read_committed/read_uncommitted). Вся эта механика подробно разобрана на живых логах ребаланса в статье про consumer groups.
  • Transactional producer: реализует протокол initTransactions/beginTransaction/commit/abort и атомарный sendOffsetsToTransaction для exactly-once — механика разобрана в статье про транзакции, включая честную асимметрию API между franz-go и kafka-clients.

Всё перечисленное — код клиентской библиотеки, а не брокера. Брокер лишь исполняет протокольные запросы (Produce, Fetch, JoinGroup, OffsetCommit, транзакционные RPC), которые формирует клиент. Поэтому «клиент Kafka» — это не деталь подключения, а самостоятельный слой логики, и от того, какой конкретно клиент выбран, зависит, что из перечисленного выше вообще доступно и как оно себя поведёт.

Три родословные

За экосистемой из десятков клиентских библиотек на самом деле стоят всего три родословные — и именно родословная, а не язык, определяет, чего от клиента ожидать.

Референсный JVM-клиентorg.apache.kafka:kafka-clients. Он же и есть часть проекта Apache Kafka: новые KIP-фичи появляются здесь первыми, а все остальные клиенты в лучшем случае их догоняют. Scala не имеет собственного протокольного клиента — fs2-kafka, zio-kafka, Alpakka Kafka оборачивают тот же kafka-clients функциональным API поверх.

Обёртки над librdkafka — библиотекой на чистом C, которая де-факто стала общим движком для всех языков, где нет своего JVM. Python (confluent-kafka), C# (Confluent.Kafka), Rust (rdkafka), Go (confluent-kafka-go), C++ (modern-cpp-kafka поверх librdkafka) — все пять вызывают одну и ту же C-библиотеку через FFI/CGo/P-Invoke. Богатый и проверенный движок, но тянет за собой нативную зависимость со всеми последствиями для сборки (ниже — отдельный раздел).

Чистые реализации протокола — библиотеки, которые сами говорят по проводу протоколом Kafka на языке хоста, без обращения ни к JVM, ни к C. На Go таких сразу три: franz-go, IBM/sarama, segmentio/kafka-go. На Python — kafka-python-ng. Плюс от отсутствия внешней зависимости выигрывает сборка и cross-compile, а минус — паритет фич обычно отстаёт от двух первых родословных, и отстаёт неравномерно: конкретно franz-go в этом смысле исключение, о котором ниже.

Таблица ниже — сводка по всем 11 драйверам, реально собранным и прогнанным живьём против одного и того же трёхброкерного KRaft-кластера (apache/kafka:4.3.1) с одинаковым сценарием: producer с ключом и (где есть ручка) идемпотентностью отправляет 12 сообщений (4 ключа order-1order-4 × 3 раунда) в топик с 3 партициями и RF=3, а consumer в группе с ручным коммитом читает их обратно.

Язык Драйвер Родословная Нативная зависимость EOS/транзакции Cooperative-sticky Produce-режим
Go franz-go v1.21.5 чистая нет да, проверено живьём да, дефолт (CooperativeStickyBalancer), проверено живьём sync (ProduceSync)
Go sarama (IBM) v1.46.1 чистая нет да, по докам (в этом сценарии не гонялось) есть, по докам (в этом сценарии не гонялось) callback (ConsumerGroupHandler)
Go segmentio/kafka-go v0.4.49 чистая нет нет — подтверждено по исходнику нет — подтверждено по исходнику sync (WriteMessages)
Go confluent-kafka-go/v2 v2.12.0 обёртка librdkafka да (CGo) да, по докам да, проверено живьём delivery-канал
Java kafka-clients референс JVM нет да, эталон да, проверено живьём sync (send().get())
Scala fs2-kafka 3.6.0 JVM-обёртка (kafka-clients 3.8.1 внутри) нет библиотека умеет (TransactionalKafkaProducer), в сценарии не гонялось наследует стратегии kafka-clients (отдельно не проверялось) функциональные потоки (fs2.Stream)
Python confluent-kafka 2.6.1 обёртка librdkafka да, но как wheel (manylinux) да, по докам да, проверено живьём callback (on_delivery)
Python kafka-python-ng 2.2.3 чистая нет частично — ручки enable.idempotence нет нет (eager sticky KIP-54, не incremental KIP-429) sync (future.get())
Rust rdkafka 0.36 обёртка librdkafka да, опционально (bundled-сборка) да, по докам да, проверено живьём async/await (FutureProducer)
C# Confluent.Kafka 2.6.1 обёртка librdkafka да, транзитивно через NuGet да, по докам да, проверено живьём async/await (ProduceAsync)
C++ modern-cpp-kafka + librdkafka обёртка librdkafka да да, по докам да, проверено живьём (cooperative[enabled] в логе) callback (delivery-лямбда)

Все 11 клиентов на живом прогоне дали одинаковый итог — 12 отправлено, 12 получено, ни одной ошибки. То, где они расходятся, — не в этом простом сценарии, а в деталях, о которых честно сказано в таблице пометками «по докам, не гонялось»: где команда стенда сознательно не дублировала проверку, потому что результат ожидался тем же (тот же принцип экономии усилий, что и в остальных статьях серии), это отмечено явно, а не выдано за проверенный факт.

Один ключ — три партиции: живой прогон

Вот флагманский факт этой статьи, и он детерминированный — воспроизводится байт в байт на любой машине с теми же ключами и тем же числом партиций, в отличие от throughput-чисел в других статьях серии. Один и тот же набор ключей order-1order-4 на одном и том же топике (3 партиции, RF=3) даёт три разных, но каждое внутренне консистентное, распределение по партициям — потому что три родословные считают партицию по ключу тремя разными формулами.

JVM-семья (murmur2 % partitions)order-1→1, order-2→0, order-3→0, order-4→2, то есть распределение 1, 0, 0, 2. Сюда попадают kafka-clients (эталон), franz-go (UniformBytesPartitioner повторяет ту же формулу (murmur2(key) & 0x7fffffff) % partitions, что и DefaultPartitioner в Java — прямая ссылка в исходнике franz-go на класс kafka-clients), kafka-python-ng, и segmentio/kafka-go — но только если явно поставить Murmur2Balancer{}, что и сделано в этом стенде: собственный балансировщик по умолчанию у segmentio даёт иное распределение.

Семья librdkafkaorder-1→1, order-2→0, order-3→2, order-4→2, то есть 1, 0, 2, 2. Здесь та же партиция для order-1 и order-2, что и в JVM-семье, но order-3 и order-4 оба уходят в партицию 2 — librdkafka по умолчанию использует собственную вариацию CRC32-хеша (consistent_random), отличную от Java DefaultPartitioner. Сюда попадают Python confluent-kafka, Rust rdkafka, C# Confluent.Kafka, C++ modern-cpp-kafka и Go confluent-kafka-go — пять независимых языковых обёрток, и все пять дали побайтово идентичное распределение. Это не совпадение, а прямое следствие того, что все пять вызывают один и тот же C-движок через разные FFI-механизмы — не пять реализаций партиционера, а одна.

sarama — третий, ни на что не похожий вариант: order-1→1, order-2→1, order-3→0, order-4→0, то есть 1, 1, 0, 0. Собственный хешер на основе FNV1a не совпадает ни с murmur2 JVM-семьи, ни с CRC32-вариацией librdkafka.

Именно три, а не два — и это легко просчитать. Две родословные (JVM и librdkafka) дают близкие, но не идентичные раскладки, так что их несложно принять за одну; а третья (sarama) на первый взгляд выглядит просто «ещё одним хешем» на фоне первых двух. На деле все три распределения различны, и ни одно не сводится к другому.

flowchart LR KEYS["order-1..order-4
одни и те же байты ключа"] --> JVM["JVM-семья
murmur2 % partitions"] KEYS --> LRDK["librdkafka-семья
CRC32-вариация"] KEYS --> SARAMA["sarama
собственный FNV1a-хешер"] JVM --> D1["1, 0, 0, 2
kafka-clients, franz-go,
segmentio+Murmur2Balancer,
kafka-python-ng"] LRDK --> D2["1, 0, 2, 2
confluent-go, Python confluent,
Rust rdkafka, C# Confluent.Kafka,
C++ modern-cpp-kafka"] SARAMA --> D3["1, 1, 0, 0
только sarama"] style KEYS fill:#f9f3e3,stroke:#8b7355 style JVM fill:#c9e4c5,stroke:#5b8a5e style LRDK fill:#eee8d8,stroke:#8b7355 style SARAMA fill:#f4d9c6,stroke:#c67a4a style D1 fill:#f9f3e3,stroke:#5b8a5e style D2 fill:#f9f3e3,stroke:#8b7355 style D3 fill:#f9f3e3,stroke:#c67a4a

flowchart LR
  KEYS["order-1..order-4
одни и те же байты ключа"] --> JVM["JVM-семья
murmur2 % partitions"] KEYS --> LRDK["librdkafka-семья
CRC32-вариация"] KEYS --> SARAMA["sarama
собственный FNV1a-хешер"] JVM --> D1["1, 0, 0, 2
kafka-clients, franz-go,
segmentio+Murmur2Balancer,
kafka-python-ng"] LRDK --> D2["1, 0, 2, 2
confluent-go, Python confluent,
Rust rdkafka, C# Confluent.Kafka,
C++ modern-cpp-kafka"] SARAMA --> D3["1, 1, 0, 0
только sarama"] style KEYS fill:#f9f3e3,stroke:#8b7355 style JVM fill:#c9e4c5,stroke:#5b8a5e style LRDK fill:#eee8d8,stroke:#8b7355 style SARAMA fill:#f4d9c6,stroke:#c67a4a style D1 fill:#f9f3e3,stroke:#5b8a5e style D2 fill:#f9f3e3,stroke:#8b7355 style D3 fill:#f9f3e3,stroke:#c67a4a
Одни и те же ключи order-1..order-4 — три семьи партиционеров, три разных распределения
Топик demo-clients-*, 3 партиции RF=3: одни ключи — три разных распределенияJVM-семья (murmur2)1, 0, 0, 2партиция 0order-2, order-3партиция 1order-1партиция 2order-4librdkafka-семья (CRC32)1, 0, 2, 2партиция 0order-2партиция 1order-1партиция 2order-3, order-4sarama (FNV1a)1, 1, 0, 0партиция 0order-3, order-4партиция 1order-1, order-2партиция 2(пусто)11 драйверов, 7 языков, тот же кластер и тот же сценарий — распределение зависит только от родословной партиционера

Практический вывод не про экзотику, а про повседневный риск: если сервисы на разных языках пишут в один топик и порядок событий одной сущности важен (значит, ключ обязан приводить к одной и той же партиции для всех писателей), недостаточно того, что все они «используют официальный клиент Kafka для своего языка». Нужно явно проверить, какую формулу партиционирования использует конкретный драйвер, и либо унифицировать её вручную (явный Partitioner/partitioner.class, где это настраивается), либо держать за правило: все продюсеры одной сущности — на клиентах одной родословной.

Почему у Go сразу несколько драйверов

Go — единственный язык в этом стенде, где вообще нет «одного разумного выбора по умолчанию, о котором можно не думать»: четыре драйвера, три родословные, реальные различия в фичах и поведении, не только в синтаксисе.

franz-go — чистая реализация протокола, современная и по факту наиболее полная по фичам среди чистых клиентов: EOS проверен живьём (в статье про транзакцииGroupTransactSession без отдельного initTransactions()), cooperative-sticky ребаланс — единственная дефолтная стратегия (CooperativeStickyBalancer, подтверждено по pkg/kgo/config.go и проверено живьём и в стенде consumer groups, и здесь). Разумный выбор по умолчанию для нового Go-кода.

sarama (IBM) — чистая реализация, зрелая (используется годами в проде до появления franz-go), но API многословнее (callback-стиль ConsumerGroupHandler.ConsumeClaim вместо синхронного цикла), и в этом стенде транзакции/cooperative-sticky не гонялись живьём — по документации они есть (Producer.Idempotent+TransactionalID, BalanceStrategyCooperativeSticky), но это честно отмечено как непроверенное здесь.

segmentio/kafka-go — чистая реализация с самым простым API, но и самая бедная по фичам среди всех 11 драйверов стенда, причём не «по докам», а подтверждено чтением исходника v0.4.49: у Writer нет поля идемпотентности вообще (нет ручки — значит нет и гарантии), у балансировщиков group нет CooperativeStickyBalancer (только RangeGroupBalancer/RoundRobinGroupBalancer/RackAffinityGroupBalancer — KIP-429 не реализован), а WriteMessages даже не возвращает partition/offset отправленной записи, в отличие от остальных трёх Go-драйверов — заметный пробел не только в фичах, но и в observability продюсера. К этому добавляется скрытая тонкость сборки: Conn.ReadPartitions() жёстко зашивает AllowAutoTopicCreation:true в протокольный запрос, и попытка использовать его как безобидную проверку «существует ли ещё топик» тихо пересоздавала топик с одной партицией вместо трёх — раньше, чем успевал отработать явный CreateTopics. Правильное решение — не пробировать состояние вообще, а сразу ретраить CreateTopics.

confluent-kafka-go/v2 — обёртка над librdkafka через CGo: максимум фич движка (cooperative-sticky здесь тоже проверен живьём), но ценой нативной сборки — подробности в разделе про CGo ниже.

Итог: «продюсер на Go» — это не одна сущность, а как минимум четыре по-разному ведущих себя реализации, и выбор между ними определяет, доступны ли вообще идемпотентность, cooperative-ребаланс и наблюдаемость отправленной записи — не только скорость или стиль API.

Что расходится между клиентами

Паритет фич убывает по родословной: JVM-клиент получает новые KIP первым, librdkafka обычно нагоняет с заметным лагом, чистые реализации — по-разному и не всегда вообще нагоняют (segmentio/kafka-go — открытый пример из раздела выше).

EOS/транзакции: JVM (эталон) и все пять librdkafka-обёрток заявляют полную поддержку. Среди чистых — franz-go подтверждён живьём, sarama заявляет по докам, segmentio/kafka-go не поддерживает вовсе, а kafka-python-ng — честная «половинчатая» история: ручки enable.idempotence нет вообще, значит producer-часть идемпотентности недоступна в принципе, что бы ни было написано в документации о самом протоколе.

Протокол ребаланса: cooperative-sticky (инкрементальный отзыв только переезжающих партиций) против eager (полный stop-the-world отзыв у всех при любом изменении состава группы) — детально с живыми логами разобрано в статье про consumer groups. Здесь важно, что дефолт различается даже внутри одной родословной: franz-go использует cooperative-sticky по умолчанию, а classic-протокол kafka-clients фактически ведёт себя как eager, если явно не сконфигурирован CooperativeStickyAssignor. Пять librdkafka-обёрток в этом стенде явно настроены на cooperative-sticky и это подтверждено живьём — в том числе прямой строкой в логе librdkafka у C++-клиента: cooperative[enabled]. А вот kafka-python-ng — отдельный случай: его sticky-ассайнер подтверждён по исходнику как классический eager sticky (KIP-54), а не incremental cooperative (KIP-429) — у объекта нет атрибута rebalance_protocol, который был бы обязателен для настоящего cooperative-режима.

Сжатие тоже не одинаково по умолчанию между клиентами: например, franz-go без явного NoCompression() по умолчанию сжимает батчи алгоритмом Snappy — деталь, которая всплыла именно как побочный эффект отладки другого сценария (реальный размер лога на диске оказался почти на порядок меньше расчётного) в статье про retention и компрессию. Если сравнивать «размер на диске» между клиентами напрямую, не зафиксировав явно кодек с обеих сторон, легко принять разницу в дефолтах клиента за разницу в эффективности алгоритма.

SASL/OAuth и безопасность — тема, где полнота покрытия у разных клиентов исторически неравномерна: SCRAM, OAUTHBEARER, mTLS поддерживаются везде из таблицы выше, но версии протоколов и удобство настройки (особенно для OAUTHBEARER с кастомным token-провайдером) отличаются заметно между родословными — librdkafka и kafka-clients обычно наиболее полны, у части чистых клиентов часть механизмов появляется позже или требует больше самостоятельной интеграции.

CGo и нативная зависимость против чистого клиента

Родословная клиента — это ещё и родословная сборки, и разница между ними не абстрактная.

Чистые клиенты — один бинарь, никакой нативной зависимости, тривиальный cross-compile: franz-go, sarama, segmentio/kafka-go в стенде запускались одинаково просто — go run . в стандартном образе golang, без единого дополнительного системного пакета.

Обёртки над librdkafka устроены иначе, и живые находки стенда это подтверждают буквально:

docker run --rm --network kafka-cookbook-net -v "$(pwd)/clients/go/confluent:/app" -w /app golang:1.25 \
  sh -c "apt-get update -qq && apt-get install -y -qq librdkafka-dev pkg-config >/dev/null && \
         CGO_ENABLED=1 go run . -brokers=kafka1:9092,kafka2:9092,kafka3:9092"

Ещё нагляднее — реальная поломка сборки на Rust. Пакет rdkafka-sys 4.10.0 (тянется крейтом rdkafka 0.36) при dynamic-linking требует системную librdkafka >= 2.12.1, а apt install librdkafka-dev в образе rust:1 (Debian trixie) даёт версию 2.8.0 — pkg-config падает несовпадением версий прямо на этапе сборки, не в рантайме. Рабочий вариант — вообще не линковаться с системной библиотекой, а собрать librdkafka из bundled C-исходников силами самого крейта:

docker run --rm --network kafka-cookbook-net -v "$(pwd)/clients/rust:/app" -w /app rust:1 \
  sh -c "apt-get update -qq && apt-get install -y -qq cmake perl >/dev/null && \
         cargo run --release -- kafka1:9092,kafka2:9092,kafka3:9092"

Здесь не нужен librdkafka-dev вообще — нужны только cmake и perl, чтобы крейт скомпилировал C-библиотеку с нуля во время сборки. Дольше собирается, зато не зависит от того, какая версия librdkafka лежит в апстрим-репозитории конкретного дистрибутива на момент сборки образа — а это ровно та проблема, которая и сломала dynamic-linking.

Практическое следствие для деплоя — то же, что при выборе базового образа для multi-stage сборок: чистый клиент кросс-компилируется в минимальный scratch/distroless-образ без сложностей, а librdkafka-обёртка требует либо совместимого runtime-образа с нужной версией .so, либо статической линковки нативной библиотеки в билд-стадии — дополнительный шаг, который стоит закладывать в CI заранее, а не находить постфактум, как это произошло с Rust-стендом.

Живой код: produce и consume на трёх языках

Сценарий одинаков на всех 11 драйверах: producer с ключом и (где есть ручка) идемпотентностью шлёт 12 сообщений (4 ключа order-1order-4 × 3 раунда) в топик с 3 партициями, RF=3 — consumer в группе с явно выключенным автокоммитом читает их обратно и печатает (partition, offset, key). Ниже — три контрастных языка: чистая реализация на Go, референсный JVM-клиент, librdkafka-обёртка на Rust.

// idempotence в franz-go включена ПО УМОЛЧАНИЮ (нет отдельного флага
// "enable.idempotence" — она живёт в kgo.Client всегда, если явно не
// отключена kgo.DisableIdempotentWrite()). acks=all — для наглядности явно.
func produce(seeds []string) int {
    cl, err := kgo.NewClient(
        kgo.SeedBrokers(seeds...),
        kgo.RequiredAcks(kgo.AllISRAcks()),
    )
    if err != nil {
        log.Fatalf("kgo.NewClient (producer): %v", err)
    }
    defer cl.Close()

    ctx := context.Background()
    count := 0
    for round := 0; round < 3; round++ {
        for _, k := range keys {
            value := fmt.Sprintf("%s-evt-%d", k, round)
            rec := &kgo.Record{Topic: topic, Key: []byte(k), Value: []byte(value)}
            res, err := cl.ProduceSync(ctx, rec).First()
            if err != nil {
                log.Fatalf("produce key=%s: %v", k, err)
            }
            fmt.Printf("  sent  key=%s partition=%d offset=%d\n", k, res.Partition, res.Offset)
            count++
        }
    }
    return count
}
private static int produce(String brokers) {
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
    props.put(ProducerConfig.ACKS_CONFIG, "all");
    props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);

    int count = 0;
    try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
        for (int round = 0; round < 3; round++) {
            for (String key : KEYS) {
                ProducerRecord<String, String> record =
                        new ProducerRecord<>(TOPIC, key, key + "-evt-" + round);
                RecordMetadata md = producer.send(record).get(); // синхронно, как ProduceSync
                System.out.printf("  sent  key=%s partition=%d offset=%d%n",
                        key, md.partition(), md.offset());
                count++;
            }
        }
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
    return count;
}
async fn consume(brokers: &str, expected: usize) -> Vec<String> {
    let consumer: StreamConsumer = ClientConfig::new()
        .set("bootstrap.servers", brokers)
        .set("group.id", GROUP_ID)
        .set("auto.offset.reset", "earliest")
        .set("enable.auto.commit", "false")
        // явный cooperative-sticky — librdkafka по умолчанию его НЕ включает
        .set("partition.assignment.strategy", "cooperative-sticky")
        .create()
        .expect("StreamConsumer::create");

    consumer.subscribe(&[TOPIC]).expect("subscribe");

    let mut recs: BTreeMap<(i32, i64), String> = BTreeMap::new();
    while recs.len() < expected {
        match consumer.recv().await {
            Ok(msg) => {
                let key = msg.key().map(|k| String::from_utf8_lossy(k).to_string()).unwrap_or_default();
                recs.insert((msg.partition(), msg.offset()), key);
                // ручной синхронный коммит offset после КАЖДОЙ записи
                consumer.commit_message(&msg, CommitMode::Sync).expect("commit_message");
            }
            Err(e) => panic!("consumer error: {e}"),
        }
    }
    recs.into_iter().map(|((p, o), k)| format!("(partition={p}, offset={o}, key={k})")).collect()
}

Три фрагмента показывают три разных стиля API на одной и той же задаче: franz-go — синхронный ProduceSync/PollFetches без колбэков, kafka-clients — классический Future.get() поверх блокирующего сокета, rdkafka — async/await поверх tokio, где даже коммит offset остаётся синхронным вызовом внутри асинхронной функции. Ни один из стилей не «правильнее» — но переключение между ними при смене языка или драйвера требует пересмотра паттернов обработки ошибок и бэкпрешера, а не только синтаксиса.

Как выбирать

  • JVM/Scala — официальный kafka-clients, либо функциональная обёртка (fs2-kafka/zio-kafka/Alpakka) под привычный стиль. Стоит явно сверять версию kafka-clients внутри обёртки с версией брокера в проде: в этом стенде fs2-kafka 3.6.0 тянет kafka-clients 3.8.1, тогда как кластер запинен на 4.3.1 — расхождение мажорных версий клиента и брокера не сломало сценарий здесь, но это не гарантия на любой комбинации фич, и версию стоит проверять явно, а не полагаться на транзитивную зависимость.
  • Gofranz-go по умолчанию: лучший паритет фич среди чистых клиентов, cooperative-ребаланс включён из коробки, без CGo. confluent-kafka-go — если нужен весь набор возможностей librdkafka (SASL-механизмы, зрелые квоты) и нативная зависимость не проблема для деплоя. sarama — оправдан организационными причинами (уже в проде, команда знает API) скорее, чем техническими преимуществами перед franz-go. segmentio/kafka-go — для простых сценариев без идемпотентности и без потребности видеть partition/offset отправленной записи; для всего остального — явный пробел в фичах, а не вопрос вкуса.
  • Pythonconfluent-kafka для продакшна и производительности (тот же движок librdkafka, что и везде), kafka-python-ng — там, где важнее отсутствие нативной зависимости, чем полный набор гарантий producer’а.
  • Rustrdkafka фактически безальтернативен как продакшн-стандарт; закладывать bundled-сборку (cmake+perl, без dynamic-linking) в CI сразу, а не после первой поломки версии системной librdkafka.
  • C#/C++Confluent.Kafka и modern-cpp-kafka+librdkafka — тот же движок, что и везде в librdkafka-семье, зрелая и предсказуемая пара для этих экосистем.
  • Общий принцип, который стоит держать в голове независимо от языка: если несколько сервисов на разных языках пишут в общий топик и порядок в рамках ключа важен, схема партиционирования — часть контракта между сервисами, а не деталь реализации отдельного клиента. Три родословные, три формулы хеша — прежде чем полагаться на «одинаковый ключ → одна партиция» через границу языков, стоит явно проверить, что все писатели используют совместимую формулу, а не просто «официальный клиент для своего языка».

Практический выбор клиента для JVM-стека и границы с экосистемой Spring разобраны отдельно в статье «Messaging на JVM: Kafka в Spring-экосистеме». Как драйвер на Go использует горутины и каналы для конкурентного PollFetches/Produce — смежная тема статей о конкурентности в Go. Дальше в серии — финальная статья про MirrorMaker 2 и геораспределённый Kafka, где та же ISR-модель, разобранная в статье про репликацию и надёжность, масштабируется на два независимых кластера. Для навигации по всей теме messaging — карта messaging-landscape-map.

Обсуждение в Telegram

Присоединиться →

Комментарии