Привет, бойцы, вы готовы победить перенос данных?
Есть задача: перенести данные из PostgreSQL в Kafka. Первая мысль: написать небольшой сервис.
Он будет выполнять SELECT, превращать строки в сообщения и отправлять их через Kafka Producer. На схеме все выглядит почти безобидно:

Но если поближе рассмотреть данное решение сразу появляются вопросы:
как понять, какие строки уже обработаны?
где хранить позицию чтения?
что произойдет после рестарта?
Ощущаете чем попахивает? Вместо «небольшого перекладывателя данных» получается полноценный сервис. Его еще нужно развертывать, обновлять и поддерживать.
Kafka Connect закрывает большую часть этой работы готовым фреймворком. В этой статье разберемся, из каких частей он состоит, как хранит состояние, что происходит при падении воркера (worker) и где заканчиваются его возможности.
PS. Буду рад подписки на tg канал с полезной инфой — Telegram:)
Что такое Kafka Connect
Kafka Connect входит в Apache Kafka и предназначен для потокового переноса данных между Kafka и внешними системами.
У него есть два основных направления работы:

За «то и откуда» отвечает свой тип коннектора:
Source Connector читает данные из внешней системы и пишет их в Kafka
Sink Connector читает сообщения из Kafka и записывает их во внешнюю систему
Kafka Connect не знает, как читать PostgreSQL или грузить файлы в S3. Эту логику поставляют плагины. А Connect отвечает за управление жизненным циклом, хранение offsets, распределение работы, перезапуск задач и предоставления REST API.
Worker, Connector и Task
В Kafka Connect существуют три сущности
Worker
Worker по сути запущенный процесс Kafka Connect. Занимается он следующим: загружает плагины, принимает запросы по REST API, запускает коннекторы и задачи, участвует в распределении работы и взаимодействует с Kafka.
Connector
Connector описывает конкретную интеграцию.

В конфигурации указываются класс плагина, подключение к источнику, таблицы, режим чтения и остальные параметры. Connector анализирует эту конфигурацию и определяет, какие задачи нужно создать.
Task
Task это непосредственная единица работы. Именно task читает данные из источника или записывает их в приемник.
Один connector может создать несколько tasks. После этого Kafka Connect распределит их между worker'ами.
В конфигурации часто встречается параметр: tasks.max: число. Важно, что это не команда «создай четыре задачи», а верхняя граница. Плагин может создать меньше tasks, если источник нельзя распараллелить.
Откуда Connect знает, на чем остановился
Одна из главных идей Kafka Connect похожа на устройство самой Kafka.
В Kafka топик разделен на партиции, а позиция внутри каждой партиции задается offset'ом.

Kafka Connect предлагает плагину представить внешний источник как набор независимых потоков, то есть source partitions. Для каждой такой партиции хранится собственная позиция, source offset.
Например, есть таблица в PostgreSQL users с возрастающим идентификатором id. Тогда в source offset будет храниться id последней переданной записи, а в source partition будет лежать users. Если task упал, то после запуска чтение таблицы начнется примерно с такой логикой:
SELECT * FROM users WHERE id > 10 ORDER BY id
Формат source partition и source offset определяет конкретный плагин.
У sink connector'ов ситуация немного другая. Они читают обычные Kafka partitions как consumer и используют Kafka offsets потребительской группы.
Standalone и Distributed Mode
Kafka Connect можно запускать в двух режимах.
Standalone Mode
В standalone‑режиме вся работа выполняется одним процессом. Так проще начать локальную разработку или проверить конфигурацию. Source offsets в этом режиме можно хранить в локальном файле.
Distributed Mode
В distributed‑режиме несколько worker'ов образуют кластер. Для этого у них должны совпадать, среди прочего, group.id и имена внутренних топиков.
Если какой‑то worker падает, его задачи после ребалансировки перейдут на оставшиеся процессы. При добавлении нового worker'а задачи также могут быть перераспределены.
В этом режиме Kafka Connect хранит свое состояние в самой Kafka. В топиках:
connect-configs— конфигурации connector'ов и tasksconnect-offsets— offsets source connector'овconnect-statuses— состояния connector'ов и tasks
Что происходит при ошибках
По умолчанию Connect придерживается fail‑fast поведения. Ошибка при конвертации или применении SMT может завершить task.
Можно разрешить пропуск проблемных записей: errors.tolerance=all
Но сама по себе эта настройка опасна. Pipeline продолжит работать и снаружи может выглядеть здоровым, хотя часть данных уже потеряна.
Как минимум стоит включить логирование: errors.log.enable=true и errors.log.include.messages=false
SMT: небольшие изменения по дороге
Иногда сообщение нужно изменить до записи в Kafka или внешнюю систему:
удалить поле
переименовать поле
добавить статическое значение
сделать поле ключом сообщения
Для этого в Kafka Connect есть Single Message Transformations или SMT. Они применяются к каждой записи по отдельности.
Управление через REST API
Kafka Connect задуман как постоянно работающий сервис и управляется через HTTP.
Основные запросы:
GET /connector-plugins GET /connectors POST /connectors GET /connectors/{name} GET /connectors/{name}/status PUT /connectors/{name}/config PUT /connectors/{name}/pause PUT /connectors/{name}/resume PUT /connectors/{name}/stop POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true DELETE /connectors/{name}
В целом, тут достаточно интуитивный API.
Когда Kafka Connect подходит, а когда нет
Kafka Connect хорошо подходит, если задача сводится к стандартному переносу данных:
база данных — Kafka
Kafka — объектное хранилище
Kafka — поисковый индекс
Kafka — аналитическая база
Особенно он полезен, когда для обеих систем уже есть зрелый поддерживаемый connector.
Отдельный сервис разумнее, если:
внутри pipeline много доменной логики
нужны join, агрегации, окна или состояние
протокол источника нестандартный, а подходящего плагина нет
Но никто не мешает написать собственный connector.
Пример: перенос JSON‑файлов в Kafka
Нам нужно:
находить новые файлы
читать из них заказы
отправлять каждый заказ отдельным сообщением в Kafka
переименовать поле
client_idвcustomer_idдобавить к каждому сообщению поле
source
Для чтения директории будем использовать SpoolDirJsonSourceConnector. Этот плагин следит за появлением файлов, читает их и после успешной обработки переносит в отдельную директорию.
Установка плагина
Для начала поставим плагин в Docker:
FROM confluentinc/cp-kafka-connect:8.3.1 RUN confluent-hub install \ --no-prompt \ jcustenborder/kafka-connect-spooldir:latest
Проверим успешность:
curl http://localhost:8083/connector-plugins
В ответе должен присутствовать класс:

Подготавливаем директории
Создадим три директории:
mkdir -p data/input mkdir -p data/finished mkdir -p data/error
input— сюда будут поступать новые файлыfinished— сюда перемещаются успешно обработанные файлыerror— сюда попадут файлы, которые не получилось прочитать
Если Kafka Connect работает внутри Docker, директории нужно примонтировать в контейнер:
volumes: - ./data/input:/data/input - ./data/finished:/data/finished - ./data/error:/data/error
Создаем файл с заказами
Добавим в data/input файл orders-001.json:
{"order_id":1001,"client_id":42,"amount":1290}
Создаем топик
kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create \ --topic partner-orders \ --partitions 1 \ --replication-factor 1
Запускаем Source Connector
Создадим connector через REST API:
curl -X POST http://localhost:8083/connectors \ -H 'Content-Type: application/json' \ -d '{ "name": "partner-orders-source-v2", "config": { "connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirJsonSourceConnector", "tasks.max": "1", "input.path": "/data/input", "finished.path": "/data/finished", "error.path": "/data/error", "input.file.pattern": "orders-.*\\.json", "cleanup.policy": "MOVE", "halt.on.error": "true", "topic": "partner-orders", "key.schema": "{\"name\":\"partner.order.Key\",\"type\":\"STRUCT\",\"isOptional\":false,\"fieldSchemas\":{\"order_id\":{\"type\":\"INT64\",\"isOptional\":false}}}", "value.schema": "{\"name\":\"partner.order.Value\",\"type\":\"STRUCT\",\"isOptional\":false,\"fieldSchemas\":{\"order_id\":{\"type\":\"INT64\",\"isOptional\":false},\"client_id\":{\"type\":\"INT64\",\"isOptional\":false},\"amount\":{\"type\":\"INT64\",\"isOptional\":false}}}", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "transforms": "RenameClientId,AddSource", "transforms.RenameClientId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.RenameClientId.renames": "client_id:customer_id", "transforms.AddSource.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.AddSource.static.field": "source", "transforms.AddSource.static.value": "partner-file" } }'
Разберем наиболее важные параметры.
input.file.pattern задает регулярное выражение для имен файлов. В нашем случае connector будет обрабатывать файлы вроде:
orders-001.json orders-002.json orders-2026-09-03.json
После успешной обработки cleanup.policy=MOVE переместит файл из input в finished. Благодаря этому один и тот же файл не будет постоянно обрабатываться заново.
Отдельно мы подключили две SMT:
transforms=RenameClientId,AddSource
Первая трансформация переименовывает поле:
transforms.RenameClientId.renames=client_id:customer_id
Вторая добавляет новое статическое поле:
transforms.AddSource.static.field=source transforms.AddSource.static.value=partner-file
SMT выполняются в том порядке, в котором перечислены в transforms. Сначала Kafka Connect переименует поле, а затем добавит источник сообщения.
Проверяем результат
Сначала посмотрим состояние connector'а:
curl http://localhost:8083/connectors/partner-orders-source/status
Результат:

Теперь прочитаем сообщения из топика:
kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic partner-orders \ --from-beginning
В топике должны появиться запись:
{"order_id":1001,"customer_id":42,"amount":1290,"source":"partner-file"}
В исходном файле поля назывались client_id, а поля source вообще не было. Продюсер при этом мы не писали: необходимые изменения выполнил сам Kafka Connect с помощью SMT.
После обработки файл orders-001.json будет перемещен из data/input в data/finished.

