Изучаем Apache Kafka с нуля. Урок 32. connect-distributed.sh: Kafka Connect в distributed-режиме

Изучаем Apache Kafka с нуля. Урок 32. connect-distributed.sh: Kafka Connect в distributed-режиме

В уроке 31 мы разобрали standalone-режим Kafka Connect. Один процесс, конфигурация в файлах, офсеты на диске. Удобно для разработки, но в production такое не ставят: упал воркер — встали все коннекторы, и подхватить их некому.

Distributed-режим решает именно это. Несколько воркеров объединяются в группу по group.id и синхронизируются через Kafka-топики. Конфигурация, офсеты и статусы хранятся не на диске, а в брокерах. При падении одного воркера остальные перераспределяют таски между собой автоматически.

Это тот же фреймворк, та же концепция коннекторов и тасков — но уже пригодная для реальных нагрузок. Именно distributed Connect используется в большинстве production-систем, где Kafka читает из PostgreSQL, пишет в S3, синхронизирует данные с Elasticsearch.

 

Как работает Kafka Connect distributed connect-distributed.sh

Скрипт connect-distributed.sh запускает один воркер Kafka Connect и присоединяет его к группе. Никаких файлов с конфигурацией коннекторов при старте не нужно — всё управление идёт через REST API уже после запуска.

connect-distributed.sh config/connect-distributed.properties

Один файл в аргументах — это конфигурация воркера. В отличие от standalone, коннекторы сюда не передаются: их создают, останавливают и удаляют через HTTP-запросы к REST API воркера.

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Архитектура distributed-кластера

Каждый воркер в группе — самостоятельный процесс, который держит часть тасков. Когда воркер присоединяется или покидает группу, происходит ребаланс: таски перераспределяются между живыми участниками. Это чем-то похоже на ребаланс consumer group, потому что механизм буквально тот же самый.


flowchart TB
    subgraph sources ["Источники данных"]
        PG["PostgreSQL"]
        FS["File System"]
    end

    subgraph workers ["Connect Group: group.id = connect-cluster"]
        W1["Worker 1\nREST :8083\nTask A1, Task B1"]
        W2["Worker 2\nREST :8083\nTask A2"]
        W3["Worker 3\nREST :8083\nTask C1"]
    end

    subgraph kafka ["Apache Kafka"]
        T1["connect-configs\n(конфигурация)"]
        T2["connect-offsets\n(офсеты)"]
        T3["connect-status\n(статусы)"]
        TP["Топики с данными"]
    end

    subgraph sinks ["Приёмники данных"]
        S3["Amazon S3"]
        ES["Elasticsearch"]
    end

    PG --> W1
    PG --> W2
    FS --> W3
    W1 <--> T1
    W1 <--> T2
    W1 <--> T3
    W2 <--> T1
    W2 <--> T2
    W2 <--> T3
    W3 <--> T1
    W3 <--> T2
    W3 <--> T3
    W1 --> TP
    W2 --> TP
    W3 --> TP
    TP --> S3
    TP --> ES

Три служебных топика — это фундамент distributed-кластера. connect-configs хранит текущую конфигурацию всех коннекторов. connect-offsets фиксирует, до какого места дочитал каждый source-коннектор. connect-status содержит состояние воркеров и тасков. Kafka при этом выступает не только транспортом данных, но и распределённым хранилищем состояния всего кластера Connect.

 

Шаг 1. Конфигурация воркера

Готовый шаблон лежит в config/connect-distributed.properties. Разберём ключевые параметры, которые нужно настроить перед запуском.

# /opt/kafka/config/connect-distributed.properties

# К каким брокерам подключается воркер
bootstrap.servers=localhost:9092

# Уникальный идентификатор группы воркеров — все воркеры с одинаковым group.id
# образуют один кластер Connect
group.id=connect-cluster

# Сериализация ключей и значений по умолчанию
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true

# Служебные топики для хранения состояния кластера
# Создайте их заранее с нужным replication.factor
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status

# Фактор репликации для служебных топиков
# В production ставим не меньше 3
config.storage.replication.factor=3
offset.storage.replication.factor=3
status.storage.replication.factor=3

# REST API: порт и хост воркера
rest.port=8083
rest.host.name=0.0.0.0

# Путь к плагинам — коннекторам, трансформациям, конвертерам
plugin.path=/opt/kafka/plugins,/usr/share/java

Параметр group.id — критичный. Все воркеры с одинаковым group.id работают как один кластер и делят между собой таски. Если запустить воркер с другим group.id рядом с тем же брокером — это будет отдельный изолированный кластер Connect.

 

Шаг 2. Создание служебных топиков

Перед первым запуском создайте три служебных топика вручную. Connect создаст их сам при старте, но с настройками по умолчанию — а это часто не то, что нужно в production. Лучше контролировать параметры явно.

# connect-configs: хранит конфигурации коннекторов
# cleanup.policy=compact обязателен — Connect дочитывает историю при рестарте
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic connect-configs \
  --partitions 1 \
  --replication-factor 3 \
  --config cleanup.policy=compact

# connect-offsets: хранит офсеты source-коннекторов
# Много партиций — выше параллелизм при большом числе коннекторов
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic connect-offsets \
  --partitions 25 \
  --replication-factor 3 \
  --config cleanup.policy=compact

# connect-status: статусы воркеров и тасков
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic connect-status \
  --partitions 5 \
  --replication-factor 3 \
  --config cleanup.policy=compact

Для всех трёх топиков обязателен cleanup.policy=compact. Connect при старте воркера перечитывает эти топики с начала, чтобы восстановить состояние. Если политика будет delete — часть истории удалится, и воркер потеряет контекст.

 

Шаг 3. Запуск воркеров

Запустить первый воркер просто — одна команда с файлом конфигурации.

# Запуск в foreground (для отладки)
connect-distributed.sh config/connect-distributed.properties

# Запуск в фоне с записью логов
connect-distributed.sh config/connect-distributed.properties &> /var/log/kafka/connect.log &

# Проверка: воркер поднял REST API
curl -s http://localhost:8083/ | python3 -m json.tool

Для второго и третьего воркеров конфигурация идентичная — достаточно скопировать файл на другие машины и запустить там же команду. Воркеры найдут друг друга через group.id и брокеры Kafka. Никаких дополнительных настроек не нужно.

Убедиться, что все воркеры в группе подключились, можно через REST API любого из них.

# Список всех участников кластера Connect
curl -s http://localhost:8083/connectors | python3 -m json.tool

# Информация о самом воркере
curl -s http://localhost:8083/ 

 

Шаг 4. Управление коннекторами через REST API

Вся работа с distributed Connect происходит через HTTP. Создать коннектор, узнать его статус, перезапустить таск, удалить — всё это POST/GET/DELETE-запросы к порту 8083. Файлы с конфигурацией коннекторов уже не нужны.

# Создать source-коннектор (FileStreamSource для теста)
curl -s -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "file-source",
    "config": {
      "connector.class": "FileStreamSource",
      "tasks.max": "1",
      "file": "/tmp/test-input.txt",
      "topic": "connect-file-test"
    }
  }' | python3 -m json.tool

# Проверить статус коннектора
curl -s http://localhost:8083/connectors/file-source/status | python3 -m json.tool

# Создать sink-коннектор (FileStreamSink)
curl -s -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "file-sink",
    "config": {
      "connector.class": "FileStreamSink",
      "tasks.max": "1",
      "file": "/tmp/test-output.txt",
      "topics": "connect-file-test"
    }
  }' | python3 -m json.tool

# Список всех коннекторов
curl -s http://localhost:8083/connectors

# Статус конкретного таска
curl -s http://localhost:8083/connectors/file-source/tasks/0/status

# Перезапустить таск (если завис или упал)
curl -s -X POST http://localhost:8083/connectors/file-source/tasks/0/restart

# Остановить коннектор (без удаления)
curl -s -X PUT http://localhost:8083/connectors/file-source/pause

# Возобновить коннектор
curl -s -X PUT http://localhost:8083/connectors/file-source/resume

# Удалить коннектор
curl -s -X DELETE http://localhost:8083/connectors/file-source

Запросы можно отправлять к любому воркеру в группе — они синхронизируются через Kafka. Создали коннектор через воркер на узле 1, статус видно через воркер на узле 2 и 3.

 

Шаг 5. Обновление конфигурации без перезапуска

Одно из главных преимуществ distributed-режима — конфигурацию коннектора можно обновить на лету, через PUT-запрос. Воркеры применят изменения и перераспределят таски без остановки всей группы.

# Проверено: Apache Kafka 4.2.0, Ubuntu 22.04

# Обновить конфигурацию существующего коннектора
# PUT заменяет конфигурацию целиком — передавайте все параметры, не только изменённые
curl -s -X PUT http://localhost:8083/connectors/file-source/config \
  -H "Content-Type: application/json" \
  -d '{
    "connector.class": "FileStreamSource",
    "tasks.max": "2",
    "file": "/tmp/test-input.txt",
    "topic": "connect-file-test"
  }' | python3 -m json.tool

Параметр tasks.max в этом примере вырос с 1 до 2. Connect сделает ребаланс тасков между воркерами автоматически — примерно то же, что происходит с partitions в consumer group при изменении числа участников.

 

Ребаланс и отказоустойчивость

Самое важное отличие distributed от standalone — поведение при сбое. Когда воркер в группе падает, остальные замечают это через heartbeat-механизм и запускают ребаланс. Таски упавшего воркера переходят к живым участникам группы.

# Смотрим на статус коннектора во время ребаланса
watch -n 2 'curl -s http://localhost:8083/connectors/file-source/status | python3 -m json.tool'

# Статус тасков после ребаланса — таски должны быть RUNNING
curl -s http://localhost:8083/connectors/file-source/status

# Проверить, на каком воркере сейчас работает таск
curl -s http://localhost:8083/connectors/file-source/tasks/0/status

Ребаланс занимает несколько секунд — ровно столько, сколько нужно воркерам на согласование через connect-status топик. В это время таски временно приостановлены. Source-коннекторы при возобновлении читают данные с того места, где остановились — офсеты в connect-offsets сохраняются.

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Distributed vs Standalone

Оба режима запускают один и тот же фреймворк с одними и теми же коннекторами. Разница — в том, где хранится состояние и как управляется жизненный цикл.

  • Standalone. Конфигурация в файлах, офсеты на диске, управление через аргументы при запуске. Подходит для разработки, одноразовых миграций и случаев, когда отказоустойчивость не нужна.
  • Distributed. Конфигурация и офсеты в Kafka, управление через REST API, автоматический ребаланс при падении воркера. Это production-вариант для всех задач интеграции, где важна непрерывность.

Переходить между режимами нельзя на лету: если коннектор работал в standalone, перенести его в distributed означает создать заново через REST API и заново указать, с какой позиции читать.

Полный цикл развёртывания production-кластера Kafka Connect с мониторингом через JMX и Prometheus, настройкой JDBC Connector и S3 Sink Connector разбирается в курсе «Apache Kafka: Администрирование кластера». Там же практика ребаланса на живом кластере и разбор типичных проблем с коннекторами.

 

Что дальше

Следующий урок — финальный. В уроке 33 разберём kcat — альтернативный CLI-инструмент для работы с Kafka, который умеет читать, писать и инспектировать кластер гораздо быстрее и удобнее, чем стандартные скрипты из bin/.

Референсы

Все уроки курса

Тема URL
1 Установка Kafka с Zookeeper https://test.bigdataschool.ru/blog/news/lesson1-kafka-zookeeper-install/
2 Установка Kafka в режиме KRaft https://test.bigdataschool.ru/blog/news/lesson2-kafka-kraft-install/
3 Docker KRaft: однонодовый кластер https://test.bigdataschool.ru/blog/news/lesson3-kafka-docker-single/
4 Docker KRaft: 3-нодовый кластер https://test.bigdataschool.ru/blog/news/lesson4-kafka-docker-cluster/
5 Утилиты bin/: переменные окружения и основы https://test.bigdataschool.ru/blog/news/lesson5-kafka-bin-intro/
6 kafka-topics.sh: управление топиками https://test.bigdataschool.ru/blog/news/lesson6-kafka-topics/
7 kafka-console-producer.sh https://test.bigdataschool.ru/blog/news/lesson7-kafka-console-producer/
8 kafka-console-consumer.sh https://test.bigdataschool.ru/blog/news/lesson8-kafka-console-consumer/
9 kafka-server-start.sh / kafka-server-stop.sh https://test.bigdataschool.ru/blog/news/lesson9-kafka-server-start-stop/
10 kafka-storage.sh https://test.bigdataschool.ru/blog/news/lesson10-kafka-storage/
11 Шпаргалка по кластеру Kafka https://test.bigdataschool.ru/blog/news/lesson11-kafka-cluster-cheatsheet/
12 kafka-configs.sh (broker) https://test.bigdataschool.ru/blog/news/lesson12-kafka-configs-broker/
13 kafka-metadata-shell.sh https://test.bigdataschool.ru/blog/news/lesson13-kafka-metadata-shell/
14 kafka-features.sh https://test.bigdataschool.ru/blog/news/lesson14-kafka-features/
15 kafka-configs.sh (topics & clients) https://test.bigdataschool.ru/blog/news/lesson15-kafka-configs/
16 kafka-log-dirs.sh https://test.bigdataschool.ru/blog/news/lesson16-kafka-log-dirs/
17 kafka-dump-log.sh https://test.bigdataschool.ru/blog/news/lesson17-kafka-dump-log/
18 kafka-delete-records.sh https://test.bigdataschool.ru/blog/news/lesson18-kafka-delete-records/
19 kafka-consumer-groups.sh https://test.bigdataschool.ru/blog/news/lesson19-kafka-consumer-groups/
20 kafka-streams-application-reset.sh https://test.bigdataschool.ru/blog/news/lesson20-kafka-streams-reset/
21 kafka-leader-election.sh https://test.bigdataschool.ru/blog/news/lesson21-kafka-leader-election/
22 kafka-reassign-partitions.sh https://test.bigdataschool.ru/blog/news/lesson22-kafka-reassign-partitions/
23 kafka-replica-verification.sh https://test.bigdataschool.ru/blog/news/lesson23-kafka-replica-verification/
24 kafka-acls.sh https://test.bigdataschool.ru/blog/news/lesson24-kafka-acls/
25 kafka-broker-api-versions.sh https://test.bigdataschool.ru/blog/news/lesson25-kafka-broker-api-versions/
26 kafka-get-offsets.sh https://test.bigdataschool.ru/blog/news/lesson26-kafka-get-offsets/
27 kafka-verifiable-producer/consumer.sh https://test.bigdataschool.ru/blog/news/lesson27-kafka-verifiable/
28 kafka-producer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson28-kafka-producer-perf/
29 kafka-consumer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson29-kafka-consumer-perf/
30 kafka-mirror-maker.sh https://test.bigdataschool.ru/blog/news/lesson30-kafka-mirror-maker/
31 connect-standalone.sh https://test.bigdataschool.ru/blog/news/lesson31-kafka-connect-standalone/
32 connect-distributed.sh https://test.bigdataschool.ru/blog/news/lesson32-kafka-connect-distributed/
33 kcat. Альтернативный CLI https://test.bigdataschool.ru/blog/news/lesson33-kcat/