Содержание
- Как работает Kafka Connect distributed connect-distributed.sh
- Архитектура distributed-кластера
- Шаг 1. Конфигурация воркера
- Шаг 2. Создание служебных топиков
- Шаг 3. Запуск воркеров
- Шаг 4. Управление коннекторами через REST API
- Шаг 5. Обновление конфигурации без перезапуска
- Ребаланс и отказоустойчивость
- Distributed vs Standalone
- Что дальше
- Референсы
- Все уроки курса
В уроке 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/.
Референсы
- Apache Kafka Connect Documentation — kafka.apache.org (2025)
- Kafka Connect Worker Configuration Reference — kafka.apache.org (2025)
- Kafka Connect REST API Reference — kafka.apache.org (2025)
- Confluent Kafka Connect User Guide — docs.confluent.io (2026)
- KIP-980: Incremental Cooperative Rebalancing for Kafka Connect — cwiki.apache.org (2025)
