Привет, бойцы, вы готовы победить перенос данных?

Есть задача: перенести данные из 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'ов и tasks

  • connect-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

Нам нужно:

  1. находить новые файлы

  2. читать из них заказы

  3. отправлять каждый заказ отдельным сообщением в Kafka

  4. переименовать поле client_id в customer_id

  5. добавить к каждому сообщению поле 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.

Ссылки

Tg канал, где много полезной инфы

Telegram