Kafka редко живёт в одиночку — вокруг него выросла экосистема, которая закрывает типовые задачи без единой строки своего кода: Connect тянет данные из файлов, БД, очередей в топики и обратно, Schema Registry следит, чтобы producer и consumer не разошлись в формате данных, ksqlDB считает потоковые агрегаты на SQL. Знать эту экосистему полезно хотя бы для того, чтобы не переписывать руками то, что уже решено, — и наоборот, понимать, где готовый инструмент не подходит и нужен свой код.
Это седьмая, финальная по ядру статья серии «Kafka: глубокое погружение» (дальше — клиенты в разных языках и геораспределение). Опирается на модель лога и партиций из первой статьи и на идемпотентный producer из статьи про exactly-once и транзакции — Connect и его коннекторы точно так же producer и consumer, только чужие, написанные и упакованные за вас. Все факты ниже — из живого прогона Schema Registry, Kafka Connect и сериализации на трёхброкерном KRaft-кластере серии, включая честную несостыковку с ожиданием и то, что не разворачивалось вовсе. Исходники всего стенда — Connect, Schema Registry, сериализация — в каталоге kafka/ecosystem публичного репозитория digital-cookbook.
В статье
- Kafka Connect: интеграция без своего кода
- Schema Registry: почему Apicurio, а не Confluent
- Эволюция схемы живьём: v1 → v2 → v3
- Совместимость: что именно проверяют BACKWARD и FORWARD
- Сериализация: цена JSON против Avro и Protobuf
- ksqlDB и streaming SQL — обзорно
- Что брать из экосистемы, а что писать самому
Kafka Connect: интеграция без своего кода
Kafka Connect решает узкую, но частую задачу: перенести данные между Kafka и внешней системой (файл, БД, очередь, объектное хранилище), не пиша для этого producer или consumer вручную. Source-коннектор читает из внешней системы и пишет в топик, sink-коннектор читает топик и пишет во внешнюю систему. И тот и другой — обычный JVM-код, упакованный в плагин, который просто конфигурируется JSON’ом через REST API worker’а, а не пишется заново под каждую интеграцию.
На практике Connect почти всегда запускают в distributed-режиме: несколько процессов worker’а образуют кластер, разделяют между собой задачи (tasks) коннекторов и хранят три вида состояния во внутренних топиках самой Kafka — конфигурацию коннекторов (config.storage.topic), позицию чтения источника (offset.storage.topic) и статус запуска (status.storage.topic). Ни отдельной БД, ни файловой системы worker’у для координации не нужно: он использует ту же Kafka, данные которой перекачивает.
bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092
group.id=ecosystem-connect-cluster
# Внутренние топики Connect — RF=3, совпадает с дефолтом кластера
offset.storage.topic=connect-offsets
offset.storage.replication.factor=3
config.storage.topic=connect-configs
config.storage.replication.factor=3
status.storage.topic=connect-status
status.storage.replication.factor=3
offset.flush.interval.ms=5000На живом кластере worker поднят на том же образе apache/kafka:4.3.1, что и брокеры (никакого отдельного Confluent-дистрибутива для запуска Connect не потребовалось — коннекторы FileStreamSource/FileStreamSink и SMT уже лежат в libs/ образа). Сценарий: source-коннектор читает 10 строк из input.txt, source-коннектор с двумя SMT (single message transform) в цепочке дописывает к каждой строке метаданные, sink-коннектор пишет результат обратно в файл.
{
"name": "ecosystem-file-source",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"tasks.max": "1",
"file": "/data/source/input.txt",
"topic": "ecosystem-connect-demo",
"transforms": "hoist,insertOrigin",
"transforms.hoist.type": "org.apache.kafka.connect.transforms.HoistField$Value",
"transforms.hoist.field": "line",
"transforms.insertOrigin.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.insertOrigin.static.field": "ingest_source",
"transforms.insertOrigin.static.value": "filestream-source-connector"
}
}SMT здесь решают ровно ту задачу, для которой они и придуманы, — лёгкую трансформацию «на лету», без отдельного сервиса и без кода: HoistField оборачивает голое строковое значение в структуру {"line": "..."} (примитив нельзя обогатить полем напрямую — сначала нужна структура), а InsertField добавляет к получившейся структуре статическое поле ingest_source=filestream-source-connector, чтобы в топике было видно происхождение записи. Реальный прогон: 10 строк на входе → 10 записей на выходе, каждая с добавленным ingest_source, оба коннектора (и сам connector, и его task) в состоянии RUNNING.
10 строк"] --> SRC["FileStreamSourceConnector
ecosystem-file-source"] SRC -->|"SMT: HoistField -> InsertField"| T["топик ecosystem-connect-demo"] T --> SINK["FileStreamSinkConnector
ecosystem-file-sink"] SINK --> OUT["output.txt
10 строк + ingest_source"] W["Distributed worker
connect-distributed.sh"] -.->|"offset/status/config"| OFF[("internal-топики
connect-offsets/-status/-configs")] W --> SRC W --> SINK style IN fill:#f9f3e3,stroke:#8b7355 style OUT fill:#c9e4c5,stroke:#5b8a5e style T fill:#eee8d8,stroke:#8b7355 style OFF fill:#eee8d8,stroke:#8b7355
flowchart LR
IN["input.txt
10 строк"] --> SRC["FileStreamSourceConnector
ecosystem-file-source"]
SRC -->|"SMT: HoistField -> InsertField"| T["топик ecosystem-connect-demo"]
T --> SINK["FileStreamSinkConnector
ecosystem-file-sink"]
SINK --> OUT["output.txt
10 строк + ingest_source"]
W["Distributed worker
connect-distributed.sh"] -.->|"offset/status/config"| OFF[("internal-топики
connect-offsets/-status/-configs")]
W --> SRC
W --> SINK
style IN fill:#f9f3e3,stroke:#8b7355
style OUT fill:#c9e4c5,stroke:#5b8a5e
style T fill:#eee8d8,stroke:#8b7355
style OFF fill:#eee8d8,stroke:#8b7355
Отдельного внимания заслуживает то, где именно живёт позиция чтения. FileStreamSourceConnector хранит offset (байтовое смещение в файле) во внутреннем топике connect-offsets по ключу «имя коннектора + путь файла» — и это состояние переживает удаление и пересоздание самого коннектора. Просто DELETE, а затем повторный POST того же коннектора на тот же файл не перечитывает его заново: позиция «конец файла» уже сохранена с прошлого запуска, и источник молча не отдаёт ни одной новой записи. Симметрично устроен и sink: у него это committed-offset обычной consumer-группы, имя которой совпадает с именем коннектора. Явный сброс — через REST offsets API (PUT .../stop → DELETE .../offsets → пересоздание коннектора и топика), доступный начиная с Connect 3.6 (KIP-875). Тонкость практическая: состояние коннектора живёт отдельно от его жизненного цикла, и об этом стоит помнить при переразвёртывании интеграций, а не только при первом запуске.
Schema Registry: почему Apicurio, а не Confluent
Schema Registry решает задачу, которую Connect не покрывает: договориться, в каком именно формате данные лежат в топике, и не дать producer’у случайно сломать этот формат для уже работающих consumer’ов. Сама Kafka о содержимом сообщения ничего не знает — для брокера это просто байты; форматом и его совместимостью между версиями занимается отдельный сервис-реестр, к которому клиенты обращаются при сериализации и десериализации.
Реестров несколько (Confluent Schema Registry, Apicurio Registry, AWS Glue Schema Registry и другие), и все они реализуют примерно одинаковую модель: subject (обычно топик + -key/-value), версии схемы внутри subject’а, настраиваемая политика совместимости. На живом кластере в качестве реестра поднят Apicurio Registry 3.2.6 (образ запинен на apicurio/apicurio-registry:3.2.6, режим хранения kafkasql) — и он завёлся с первой попытки против apache/kafka:4.3.1 в чистом KRaft-режиме, самостоятельно создав себе служебные топики kafkasql-journal, kafkasql-snapshots, registry-events на том же самом кластере. Отдельно Confluent cp-schema-registry не пробовался — раз Apicurio поднялся с первого раза и покрывает нужный REST-контракт (ccompat, Confluent-совместимый API /apis/ccompat/v7/..., теми же путями и кодами ответов), дублировать попытку ради того же самого результата смысла не было. Выбор именно Apicurio для этого стенда объясняется практично: он не требует Confluent Platform целиком, хранит свои данные в топиках того же брокера (никакого отдельного хранилища для реестра) и говорит тем же REST-диалектом, которым пользуется большинство клиентских библиотек, написанных под Confluent Schema Registry.
Здесь же — самая поучительная находка всего стенда, и она не про Kafka как таковую, а про то, как безобидная настройка брокера ломает сторонний сервис. Apicurio перед стартом верифицирует конфигурацию своих топиков и отказывается подниматься, если хранение в них не бесконечное:
io.apicurio.registry.exception.RuntimeAssertionFailedException: Runtime assertion failed:
To prevent accidental loss of data, Apicurio Registry verifies that Kafka topics are
configured correctly before starting. The following issue was found with topic
'kafkasql-journal': Topic must have 'retention.ms=-1' and 'retention.bytes=-1' to prevent
accidental loss of data. Effective configuration value is 'retention.ms=604800000' and
comes from built-in default configuration.Проверка разумная: журнал KafkaSQL — это и есть база реестра, и недельный retention означал бы, что схемы старше недели однажды просто исчезнут. Загвоздка в том, кто создаёт топик. Сам реестр создаёт его правильно — одна партиция, RF=3, retention.ms=-1, retention.bytes=-1. Но внутри реестра есть гонка: consumer подписывается на kafkasql-journal раньше, чем AdminClient успевает его создать. Обращение к несуществующему топику (UnknownTopicOrPartitionException в логе) при включённом на брокере auto.create.topics.enable=true — а это дефолт Kafka — заставляет брокер создать топик первым, с дефолтным недельным retention. Дальше верификация валит старт. Кто выиграет гонку, зависит от машины: на одной стенд поднимается стабильно, на другой стабильно нет — при идентичных compose-файле и версии образа.
Отключать проверку (APICURIO_KAFKASQL_TOPIC_CONFIGURATION_VERIFICATION_OVERRIDE_ENABLED=true существует) — ровно тот случай, когда «починка» состоит в снятии предохранителя. Правильный путь — не дать брокеру создать топик раньше: отдельный init-шаг создаёт все три служебных топика с нужными параметрами до старта реестра, а если они уже существуют с неверным retention, приводит их в порядок через kafka-configs --alter.
Две тонкости помельче — обе про то, что «сервис поднялся» и «сервис работает» это разные утверждения. Первая: базовый образ Apicurio (ubi10-minimal) не содержит wget, поэтому первым обходом healthcheck проверял открытый TCP-порт через bash и /dev/tcp. Такая проверка врёт: реестр показывает healthy даже при полностью погашённом кластере — порт-то слушает, — и сценарий идёт дальше к запросам, обслужить которые некому. curl в образе есть, поэтому healthcheck проверяет реальный ответ API. Вторая: без depends_on реестр стартует одновременно с брокерами и уходит в ретраи (Connection to node -1 (kafka1:9092) could not be established):
# создаёт kafkasql-journal / kafkasql-snapshots / registry-events
# с retention.ms=-1 ДО старта реестра (и чинит уже существующие)
kafkasql-topics-init:
image: apache/kafka:4.3.1
depends_on:
kafka1: { condition: service_healthy }
kafka2: { condition: service_healthy }
kafka3: { condition: service_healthy }
entrypoint: ["/bin/bash", "-c", "...kafka-topics.sh --create --if-not-exists ... --config retention.ms=-1 ..."]
schema-registry:
image: apicurio/apicurio-registry:3.2.6 # пин: плавающий latest-release ломает воспроизводимость
environment:
APICURIO_STORAGE_KIND: kafkasql
APICURIO_KAFKASQL_BOOTSTRAP_SERVERS: kafka1:9092,kafka2:9092,kafka3:9092
depends_on:
kafka1: { condition: service_healthy }
kafka2: { condition: service_healthy }
kafka3: { condition: service_healthy }
kafkasql-topics-init: { condition: service_completed_successfully }
healthcheck:
# реальный ответ API, а не просто открытый порт: с проверкой порта
# реестр показывал healthy даже при погашенном кластере
test: ["CMD-SHELL", "curl -fsS http://localhost:8080/apis/registry/v3/system/info || exit 1"]В этом же примере Connect и Schema Registry запущены раздельно: worker сконфигурирован на schemaless JsonConverter (без schemas.enable), а реестр показан отдельно через REST — в проде их обычно связывают через AvroConverter, который на лету сериализует значения по схеме из реестра, но добавление ещё одного конвертера в plugin.path — это отдельный кусок инфраструктуры (ещё один jar), не обязательный, чтобы показать, что каждый из трёх компонентов реально работает.
Эволюция схемы живьём: v1 → v2 → v3
Схема — это контракт между producer и consumer, и единственный практический вопрос — что происходит, когда контракт меняется, а часть потребителей ещё работает по старой версии. Реестр отвечает на этот вопрос не декларативно, а буквально: каждая попытка зарегистрировать новую версию схемы проверяется против политики совместимости subject’а, и несовместимая версия просто не попадает в реестр.
Subject ecosystem-user-value настроен на политику BACKWARD (явно, через PUT /config/{subject}). Три версии схемы User (Avro):
{
"type": "record",
"name": "User",
"namespace": "tech.khorost.kafka.ecosystem",
"fields": [
{ "name": "id", "type": "long" },
{ "name": "name", "type": "string" },
{ "name": "email", "type": "string" }
]
}{
"fields": [
{ "name": "id", "type": "long" },
{ "name": "name", "type": "string" },
{ "name": "email", "type": "string" },
{ "name": "age", "type": ["null", "int"], "default": null }
]
}{
"fields": [
{ "name": "id", "type": "string" },
{ "name": "name", "type": "string" },
{ "name": "email", "type": "string" },
{ "name": "age", "type": ["null", "int"], "default": null }
]
}Регистрация всех трёх версий через тот же ccompat-REST, которым пользуются реальные клиенты (упрощённо — без обвязки повторов и построения JSON-тела через python -c json.dumps, которая в оригинальном скрипте нужна была только из-за особенностей экранирования в Git Bash на Windows):
# (0) политика совместимости subject'а
curl -X PUT -H "Content-Type: application/json" \
-d '{"compatibility":"BACKWARD"}' \
"$SR/apis/ccompat/v7/config/ecosystem-user-value"
# (1) v1 — id/name/email
curl -w "HTTP_STATUS:%{http_code}\n" -X POST \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d @v1.json "$SR/apis/ccompat/v7/subjects/ecosystem-user-value/versions"
# (2) v2 — +age optional, BACKWARD-совместимо
curl -w "HTTP_STATUS:%{http_code}\n" -X POST \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d @v2.json "$SR/apis/ccompat/v7/subjects/ecosystem-user-value/versions"
# (3) v3 — id long -> string, НЕ совместимо
curl -w "HTTP_STATUS:%{http_code}\n" -X POST \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d @v3.json "$SR/apis/ccompat/v7/subjects/ecosystem-user-value/versions"Реальный результат: v1 зарегистрирована (HTTP_STATUS:200), v2 зарегистрирована (HTTP_STATUS:200 — расширение необязательным полем с дефолтом прошло проверку BACKWARD), v3 отклонена реестром (HTTP_STATUS:409) с точным объяснением причины: "reader type: STRING not compatible with writer type: LONG at /fields/0/type". После всей последовательности список версий subject’а — [1, 2]: несовместимая v3 просто не попала в историю, реестр не оставляет «сломанных» версий даже временно.
BACKWARD"} GATE -->|"200: совместимо"| REG1[("subject
ecosystem-user-value
version 1")] V2["v2: +age optional default=null"] -->|"POST /versions"| GATE GATE -->|"200: совместимо"| REG2[("version 2")] V3["v3: id long -> string"] -->|"POST /versions"| GATE GATE -->|"409: НЕсовместимо"| REJ["отклонена,
в историю версий не попала"] style GATE fill:#f9f3e3,stroke:#8b7355 style REG1 fill:#c9e4c5,stroke:#5b8a5e style REG2 fill:#c9e4c5,stroke:#5b8a5e style REJ fill:#e3b0a8,stroke:#8b4a3d
flowchart LR
V1["v1: id/name/email"] -->|"POST /versions"| GATE{"SR compat-гейт
BACKWARD"}
GATE -->|"200: совместимо"| REG1[("subject
ecosystem-user-value
version 1")]
V2["v2: +age optional default=null"] -->|"POST /versions"| GATE
GATE -->|"200: совместимо"| REG2[("version 2")]
V3["v3: id long -> string"] -->|"POST /versions"| GATE
GATE -->|"409: НЕсовместимо"| REJ["отклонена,
в историю версий не попала"]
style GATE fill:#f9f3e3,stroke:#8b7355
style REG1 fill:#c9e4c5,stroke:#5b8a5e
style REG2 fill:#c9e4c5,stroke:#5b8a5e
style REJ fill:#e3b0a8,stroke:#8b4a3d
Разбор эволюции схем как таковой, включая паттерны версионирования событий за пределами Avro-совместимости реестра, — в статье «Контракты событий и эволюция схем».
Совместимость: что именно проверяют BACKWARD и FORWARD
Название режима совместимости легко перепутать местами, потому что оно описывает не то, «в какую сторону меняется схема», а то, кто именно кого должен уметь прочитать:
BACKWARD— новая схема (та, по которой работает читатель) должна уметь прочитать данные, записанные по предыдущей схеме. Разрешено: удаление любого поля и добавление нового поля, но только optional (с default-значением) — потому что старые записи этого поля физически не содержат, и читателю нужно откуда-то взять значение по умолчанию.FORWARD— данные, записанные по новой схеме, должны уметь прочитать читатели, ещё работающие по предыдущей схеме. Разрешено: добавление любого поля (старый читатель его просто не заметит и проигнорирует) и удаление, но только optional-поля (если удалить обязательное, старый читатель, ожидающий его в данных, сломается).
Это две разные гарантии, направленные в противоположные стороны эволюции: BACKWARD защищает читателя, который обновился раньше писателей, FORWARD — читателя, который обновится позже. FULL требует одновременно и то и другое. Живой пример subject’а ecosystem-user-value демонстрирует именно BACKWARD-ветку правила: добавление optional-поля age с default=null в v2 — тот самый случай «добавление поля разрешено, если оно optional», а смена типа id в v3 не проходит ни при какой политике совместимости, потому что это не структурное изменение (появилось/исчезло поле), а изменение типа существующего — Avro не умеет читать 8-байтный long как string независимо от направления проверки.
Сериализация: цена JSON против Avro и Protobuf
Schema Registry решает вопрос совместимости, но у формата данных есть ещё одна ось — размер на проводе и на диске. Разница между самоописываемым текстовым форматом (JSON, имена полей едут вместе с каждым сообщением) и бинарным форматом со схемой, переданной сторонам отдельно (Avro, Protobuf), не абстрактная: на одних и тех же пяти записях User (id/name/email) она измеряется прямо.
// User — зеркало ecosystem/schemas/user-v1.avsc (id: long, name: string, email: string).
type User struct {
ID int64 `avro:"id" json:"id"`
Name string `avro:"name" json:"name"`
Email string `avro:"email" json:"email"`
}
// encodeProtobufUser — ручное tag-based wire-кодирование, эквивалентное
// proto3-сообщению `message User { int64 id=1; string name=2; string email=3; }`.
func encodeProtobufUser(u User) []byte {
var b []byte
b = protowire.AppendTag(b, 1, protowire.VarintType)
b = protowire.AppendVarint(b, uint64(u.ID))
b = protowire.AppendTag(b, 2, protowire.BytesType)
b = protowire.AppendString(b, u.Name)
b = protowire.AppendTag(b, 3, protowire.BytesType)
b = protowire.AppendString(b, u.Email)
return b
}
// j, _ := json.Marshal(u) // encoding/json, стандартная библиотека
// a, _ := avro.Marshal(schema, u) // github.com/hamba/avro/v2, по той же
// // схеме, что зарегистрирована в SR выше
// p := encodeProtobufUser(u) // google.golang.org/protobuf/encoding/protowireAvro закодирован по той же схеме user-v1.avsc, что зарегистрирована в Schema Registry выше — той же самой, которую в реальной интеграции клиенты получали бы из реестра, а не хранили бы у себя захардкоженной. Protobuf собран вручную через низкоуровневый protowire API (без protoc и генерации кода из .proto), но даёт побайтово те же данные на проводе, что и обычный сгенерированный Marshal().
JSON — 337 байт (базовая линия), Avro — 207 байт (0.61x, компактнее на 5 процентных пунктов, чем Protobuf в этом конкретном наборе данных), Protobuf — 222 байта (0.66x). Разница объясняется не «магией бинарного формата» вообще, а вполне конкретно: и Avro, и Protobuf не несут имена полей в каждом сообщении (схема передаётся сторонам отдельно — то самое, что делает Schema Registry), а числа кодируют компактным varint вместо текстового представления. В отличие от throughput-чисел в остальной серии, размер сериализации — величина детерминированная: она зависит только от схемы и данных, а не от хоста или загрузки сети, и воспроизводится байт-в-байт на любой машине.
Честная асимметрия: сравнение сериализации сделано только на Go — повторять тот же самый детерминированный расчёт на Java не добавило бы нового знания, кодировки одинаковы независимо от языка клиента. По той же причине Confluent Schema Registry отдельно не поднимался и не сравнивался с Apicurio — раз последний завёлся с первой попытки и покрывает нужный REST-контракт, тратить время на повторную установку ради того же результата не было смысла.
ksqlDB и streaming SQL — обзорно
ksqlDB кладёт SQL поверх модели потоков и таблиц Kafka: CREATE STREAM/CREATE TABLE объявляют, как читать топик, а SELECT ... EMIT CHANGES и материализованные представления считают агрегаты непрерывно, по мере поступления новых записей, а не разово по снапшоту данных. Для типовых задач — обогащение потока справочником, скользящие агрегаты по окну времени, простая маршрутизация по условию — это заметно быстрее написать на SQL, чем поднимать отдельный сервис на Kafka Streams.
В этом примере ksqlDB не развёрнут — тем же принципом экономии усилий, что и Cruise Control в статье про эксплуатацию: разворачивать ещё один потоковый движок ради демонстрации SQL-синтаксиса поверх уже показанной модели потоков было бы несоразмерно учебной задаче. Глубокий разбор потоковой обработки — стейт, время события против времени обработки, exactly-once на границе Kafka Streams — в статье «Kafka Streams и Flink: потоковая обработка»; ksqlDB там же упоминается как SQL-надстройка над той же моделью, которую Kafka Streams даёт программно.
Что брать из экосистемы, а что писать самому
Экосистема экономит код там, где задача типовая и хорошо описывается конфигурацией, а не логикой: перенос данных между системой и Kafka без трансформации бизнес-уровня — это Connect, а не свой producer/consumer с ручной обработкой ошибок и ретраев. Source для Debezium (CDC из БД, отдельный разбор — в статье про CDC с Debezium и PostgreSQL) устроен ровно так же, как FileStream-коннекторы выше: тот же worker, тот же REST API конфигурации, тот же жизненный цикл offset/status.
Обратная сторона: у готового инструмента есть своя операционная поверхность, и она не бесплатна. Distributed Connect worker — это ещё один процесс, три внутренних топика, своя логика ребаланса задач между worker’ами при масштабировании; offset коннектора, как показано выше, живёт отдельно от его конфигурации, и это нужно понимать, а не угадывать при инцидентах. Schema Registry добавляет сетевой вызов на каждую (де)сериализацию и ещё один компонент, который может быть недоступен именно тогда, когда нужен producer. Здесь работает простое правило: типовая интеграция без кастомной бизнес-логики — в Connect; там, где нужна нетривиальная трансформация, обогащение из нескольких источников или контроль над семантикой доставки тоньше, чем даёт SMT, — свой код, будь то простой producer/consumer или Kafka Streams.
Schema Registry стоит подключать, как только в системе больше одного независимого потребителя одного топика: без него совместимость формата данных превращается в устную договорённость между командами, которая рано или поздно нарушается тем же способом, что и v3 в примере выше, только без гарантированного HTTP 409 — прод просто перестаёт разбирать сообщения. Выбор формата (Avro/Protobuf/JSON Schema) — это выбор между компактностью и удобством отладки: JSON читаем глазами и не требует схемы для черновой интеграции, Avro и Protobuf компактнее и дают строгую проверку типов на границе, но требуют инфраструктуры вокруг схемы, показанной в этой статье.
Это закрывает интеграционный слой Kafka: как обмениваться данными с внешним миром без своего кода (Connect) и как не дать формату данных разойтись между producer и consumer (Schema Registry). Дальше в серии — клиенты Kafka в разных языках и MirrorMaker 2, гео-репликация, где Connect оказывается инфраструктурой уже не для интеграции с внешними системами, а для репликации между кластерами. Для навигации по всей теме messaging — карта messaging-landscape-map.
Комментарии