Содержание
- Что такое Kafka Connect и какие задачи он закрывает
- Архитектура Kafka Connect, worker, connector и task
- Standalone и distributed режимы
- Служебные топики конфигураций, офсетов и статусов
- Как данные проходят через конвейер Connect
- Converter и Single Message Transform
- Чем Source Connector отличается от Sink Connector
- Ограничения и подводные камни
- Когда брать Kafka Connect, а когда писать свой код
- Практика, запуск коннекторов через REST API
- Заключение
- Референсные ссылки
Kafka Connect это встроенный в Apache Kafka фреймворк интеграции, который переливает данные между кластером и внешними системами через готовые плагины-коннекторы. Source Connector забирает записи из внешней системы и кладёт их в топик, Sink Connector читает топик и отдаёт записи наружу. Разработчику при этом не нужен ни собственный Kafka Producer, ни Consumer: интеграция описывается конфигурационным файлом в формате JSON. Ретраи, офсеты, параллелизм и восстановление после сбоя уже реализованы внутри фреймворка.
Connect входит в дистрибутив Apache Kafka с версии 0.9 и остаётся штатным компонентом в актуальном релизе 4.3.0 от 22 мая 2026 года. Разберём, из чего он собран, где хранит состояние и где заканчивается зона его ответственности. Все числа и статусы ниже сняты с реального прогона на образе apache/kafka:4.3.0.
Что такое Kafka Connect и какие задачи он закрывает
По классу решений это интеграционный фреймворк с плагинной архитектурой. Сама платформа отвечает за жизненный цикл, распределение нагрузки, хранение состояния и отказоустойчивость, а логику чтения или записи приносит подключаемый коннектор. Kafka Connect работает как отдельный процесс рядом с кластером, а не внутри брокера, и общается с ним обычным клиентским протоколом.
Типовая задача такая. Есть операционная база и аналитический контур, между ними нужен поток данных, который переживает перезапуск и не теряет позицию чтения. Написать такой мост руками можно, но каждый новый источник означает новый сервис со своим деплоем, мониторингом и багами. Connect превращает эту работу в конфигурацию.
- Подключение баз данных. JDBC-коннектор тянет строки по возрастающей колонке, а CDC-платформа Debezium читает журнал изменений и отдаёт события вставки, обновления и удаления.
- Выгрузка в хранилища и поисковые движки. Elasticsearch, OpenSearch, ClickHouse, S3 и объектные хранилища Data Lake закрываются готовыми sink-плагинами.
- Файловые системы и обмен с legacy-системами. Чтение каталогов, выгрузка в HDFS, приём файлов по расписанию.
- Репликация между кластерами. MirrorMaker 2 построен поверх Connect и использует ту же механику worker и task.
Общий знаменатель у всех сценариев один: данные надо переложить из точки A в точку B, сохранив порядок и позицию чтения, и почти не трансформировать по дороге.
Архитектура Kafka Connect, worker, connector и task
В основе три сущности, и путать их не стоит, потому что каждая живёт на своём уровне.
- Worker. JVM-процесс, который исполняет плагины. Это единица развёртывания и масштабирования: добавили worker, получили больше ресурсов под задачи.
- Connector. Логическая единица конфигурации, экземпляр плагина. Сам он данные не двигает, а решает, как разбить работу на части, и порождает описания задач.
- Task. Единица параллелизма. Именно task читает из источника или пишет в приёмник, и именно task распределяются по доступным worker.
Иерархия простая: worker исполняет task, task порождены connector. Этой цепочки достаточно, чтобы предсказать поведение кластера под нагрузкой.
Standalone и distributed режимы
Standalone это один процесс, конфигурация в файлах на диске, позиция чтения в локальном файле офсетов. Режим удобен для разработки, демонстраций и сбора логов с конкретной машины. Отказоустойчивости здесь нет: упал процесс, встал поток.
Distributed это боевой вариант. Несколько worker запускаются с одинаковым group.id и объединяются в группу тем же протоколом координации, что и обычные группы потребителей. Они сами договариваются, кто какие task исполняет, а при добавлении, остановке или падении узла запускают ребалансировку и перераспределяют нагрузку по оставшимся. Конфигурация подаётся не файлом, а через REST API, который по умолчанию слушает порт 8083.
Служебные топики конфигураций, офсетов и статусов
Distributed-кластер не хранит состояние на диске worker. Вместо этого он держит его в самой Kafka, в трёх служебных топиках с политикой уплотнения. Вот что показал прогон демо-стенда, где топики создавались настройками по умолчанию.
- Топик конфигураций. На стенде это connect-configs, ровно одна партиция и cleanup.policy=compact. Здесь лежат конфигурации коннекторов и сгенерированные конфиги task, а однопартиционность обязательна, поскольку порядок изменений тут критичен.
- Топик офсетов. На стенде connect-offsets, 25 партиций. Хранит позиции чтения source-коннекторов, причём в понятном виде: запись по ключу [orders-file-source,{filename=/data/orders.txt}] со значением {position:64}.
- Топик статусов. На стенде connect-status, 5 партиций. Отсюда REST API отдаёт текущее состояние коннекторов и task.
Практическое следствие такое: worker в distributed-режиме почти не хранит состояния, и потерять узел не страшно. Потерять служебные топики, наоборот, означает потерять весь конвейер, поэтому фактор репликации у них в бою ставят не ниже трёх.
Apache Kafka для инженеров данных
Код курса
DEVKI
Ближайшая дата курса
24 августа, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Как данные проходят через конвейер Connect
Внутри фреймворка запись проходит несколько слоёв, и каждый настраивается отдельно. Схема ниже показывает полный маршрут в обе стороны.
Ключевая деталь в том, что коннектор работает с данными как со структурой, а не с байтами. Превращение структуры в байты и обратно вынесено в отдельный слой, поэтому один и тот же плагин умеет отдавать и JSON, и Avro, и просто строку.
Converter и Single Message Transform
Конвертер (converter) отвечает за сериализацию. Он симметричен по обе стороны конвейера: source-коннектор отдаёт объект, конвертер превращает его в байты для топика, sink-коннектор получает объект обратно. Штатные варианты это JsonConverter, StringConverter и ByteArrayConverter, а Avro и Protobuf приезжают из экосистемы Schema Registry. В релизе 4.2.0 у JsonConverter появился параметр schema.content по KIP-1054, который позволяет задать схему снаружи и не таскать её в каждом сообщении.
Single Message Transform (SMT) это лёгкая обработка отдельной записи на лету: переименовать поле, замаскировать колонку, выбросить лишнее, подставить топик по значению поля. Трансформации собираются в цепочку и выполняются по порядку. На стенде цепочка из двух штатных трансформаций превратила плоскую строку в структуру с дополнительным полем, то есть Struct{line=order-1001;paid;1990,src=kafka-connect-demo}. Работает SMT строго по одной записи за раз, без состояния и без окон. Как только логика требует джойна двух потоков или агрегации за интервал, SMT перестаёт годиться и задача уходит в Kafka Streams или Flink.
Чем Source Connector отличается от Sink Connector
Оба типа плагинов используют одну инфраструктуру, но ведут себя по-разному в мелочах, которые всплывают уже на эксплуатации.
| Характеристика | Source Connector | Sink Connector |
|---|---|---|
| Направление | Внешняя система в Kafka | Kafka во внешнюю систему |
| Аналог в клиентском API | Producer | Consumer |
| Где хранится позиция | Служебный топик офсетов Connect | Штатные офсеты группы потребителей |
| Чем ограничен параллелизм | Логикой плагина, например числом таблиц или файлов | Числом партиций топиков-источников |
| Типовой пример | JDBC source, Debezium, FileStreamSource | Elasticsearch sink, S3 sink, JDBC sink |
Отсюда следует прикладной вывод: у sink-коннектора нет смысла задирать tasks.max выше суммарного числа партиций, лишние task останутся без работы. На стенде это видно буквально. Коннектор с tasks.max=5 на топике из трёх партиций поднял пять task, все пять отчитались состоянием RUNNING, а партиции в группе потребителей разобрали только три из них.
Apache Kafka: администрирование кластера
Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Ограничения и подводные камни
Connect часто пытаются использовать как полноценный ETL-инструмент, и это главная ошибка при внедрении. По факту он закрывает E и L, а T остаётся снаружи. Ниже то, обо что спотыкаются на практике.
- Нет сложных преобразований. Джойны, агрегации, оконные функции и обогащение из другого потока фреймворком не поддерживаются вообще.
- Статус RUNNING не означает работу. Пример с пятью task на трёх партициях выше как раз про это. Мониторить надо не только состояние коннектора, но и распределение партиций по task, иначе простаивающая половина конвейера выглядит здоровой.
- Число task решает плагин, а не администратор. Параметр tasks.max это верхняя граница, а не требование. На стенде FileStreamSource с tasks.max=3 создал ровно один task, потому что делить один файл на части он не умеет.
- Ребалансировка стоит времени. Перезапуск worker в distributed-кластере запускает перераспределение task, и на этот момент поток останавливается. Инкрементальная кооперативная ребалансировка сгладила проблему, но не убрала её.
- Семантика доставки зависит от плагина. Exactly-once для source-коннекторов появилась в Kafka 3.3 по KIP-618, включается параметром exactly.once.source.support на каждом worker и требует поддержки со стороны конкретного коннектора. Для sink-стороны гарантия определяется идемпотентностью приёмника.
- Переопределение клиентских настроек. Конфигурация коннектора умеет менять параметры клиента, включая учётные данные. В релизе 4.2.0 по KIP-1188 добавлена политика Allowlist, где администратор явно перечисляет разрешённые к переопределению ключи.
Ни один из пунктов не отменяет пользу фреймворка, но все они говорят об одном: Connect выбирают под задачу перекладывания данных, а не под задачу их обработки.
Когда брать Kafka Connect, а когда писать свой код
Первый довод за Connect это типовая интеграция, под которую есть живой плагин. Второй довод это количество: когда интеграций много, единый способ их эксплуатировать экономит больше, чем идеально написанный сервис под одну из них. Третий довод это состав команды. Конфиг читается и правится быстрее чужого кода, поэтому конвейер могут вести инженеры данных, а не разработчики.
Свой Producer или Consumer оправдан в обратных ситуациях. Нестандартный протокол источника, бизнес-логика внутри маршрута, требование к задержке, которое не переживает ребалансировку, или преобразования сложнее SMT. Устройство коннекторов изнутри разобрано в статье Под капотом Kafka Connect: источники, приемники и коннекторы. Работу с самой платформой, включая Apache Kafka целиком, разбирают руками на курсе Apache Kafka для инженеров данных.
Практика, запуск коннекторов через REST API
По традиции весь код используемый в статье выкладываем на наш GitHub репозиторий
Демо-стенд поднимает брокер в режиме KRaft и worker Connect из одного образа. Конфигурация worker задаёт адрес кластера, идентификатор группы, конвертеры и имена служебных топиков.
# конфигурация worker, прогнано на apache/kafka:4.3.0 bootstrap.servers=kafka:9092 # общий идентификатор группы worker, по нему они находят друг друга group.id=connect-demo key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter # схема не едет внутри каждого сообщения key.converter.schemas.enable=false value.converter.schemas.enable=false # три служебных топика, здесь живёт всё состояние кластера Connect config.storage.topic=connect-configs config.storage.replication.factor=1 offset.storage.topic=connect-offsets offset.storage.partitions=25 status.storage.topic=connect-status status.storage.partitions=5 listeners=HTTP://:8083 # путь до плагинов, в поставке Kafka это jar с коннекторами FileStream plugin.path=/opt/kafka/libs/connect-file-4.3.0.jar
Сам коннектор создаётся POST-запросом на REST-эндпоинт worker. Конфигурация это JSON, где обязательны имя, класс плагина и верхняя граница числа task. Файл конфигурации из репозитория называется file_source.json.
# конфигурация source-коннектора, прогнано на apache/kafka:4.3.0
{
"name": "orders-file-source",
"config": {
// класс плагина, он должен лежать в plugin.path на worker
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
// верхняя граница параллелизма, реальное число task решает коннектор
"tasks.max": "3",
"file": "/data/orders.txt",
"topic": "demo_orders",
"value.converter": "org.apache.kafka.connect.storage.StringConverter",
"key.converter": "org.apache.kafka.connect.storage.StringConverter"
}
}
Дальше конфигурация уходит на worker, а состояние проверяется тем же API. Порт 8083 это значение по умолчанию, оно задаётся параметром listeners.
# команды прогона, Apache Kafka 4.3.0, Docker Desktop 29.4.2 # создание коннектора из файла конфигурации curl -s -X POST -H "Content-Type: application/json" \ --data @connectors/file_source.json \ http://localhost:8083/connectors # состояние коннектора и всех его task curl -s http://localhost:8083/connectors/orders-file-source/status # перезапуск упавшей задачи с номером 0 curl -s -X POST http://localhost:8083/connectors/orders-file-source/tasks/0/restart
Ответ на POST приходит сразу, а поле tasks в нём пустое: конфигурация принята, но задачи ещё не розданы. Через несколько секунд статус наполняется. Ответ на запрос статуса содержит состояние самого коннектора и каждой task отдельно. Именно на это поле вешают мониторинг: коннектор числится работающим, пока часть его task висит в состоянии FAILED.
# реальный вывод прогона, apache/kafka:4.3.0, tasks.max в конфиге равен 3
{"name":"orders-file-source",
"connector":{"state":"RUNNING","worker_id":"172.20.0.3:8083","version":"4.3.0"},
"tasks":[{"id":0,"state":"RUNNING","worker_id":"172.20.0.3:8083","version":"4.3.0"}],
"type":"source"}
Финальная проверка на стенде выглядела так. Worker перезапустили, и оба коннектора снова оказались в списке, хотя заново их никто не создавал. В файл источника дописали строку, в приёмнике стало шесть записей на шесть строк источника, без единого дубля. Конфигурация приехала из топика конфигураций, позиция чтения из топика офсетов.
Заключение
Kafka Connect закрывает узкую, но очень частую задачу: перекладывание данных между Kafka и внешними системами без написания кода под каждую пару. Плата за удобство это ограниченность преобразований и зависимость от качества стороннего плагина. В эксплуатации помогают три привычки. Держать в голове разделение на worker, connector и task. Помнить, что состояние конвейера живёт в служебных топиках Kafka. Проверять не только статус коннектора, но и реальное распределение партиций по task. С ними фреймворк ведёт себя предсказуемо.


