Содержание
- Что делает connect-standalone.sh
- Архитектура Kafka Connect
- Практическая настройка Kafka connect
- Шаг 1. Конфигурация воркера
- Шаг 2. Конфигурация source-коннектора
- Шаг 3. Конфигурация sink-коннектора
- Шаг 4. Запуск connect-standalone.sh
- Шаг 5. Проверка статуса через REST API
- Шаг 6. Проверка результата
- Управление через REST API
- Когда использовать standalone-режим
- Что дальше
- Референсы
- Все уроки курса
В уроке 30 разобрали MirrorMaker 2 — инструмент для репликации топиков между кластерами. Там мы мельком упоминали, что MM2 построен поверх Kafka Connect. Теперь пришло время разобраться с самим фреймворком.
Kafka Connect — это слой интеграции, который берёт на себя всю рутину переноса данных между Kafka и внешними системами. Читать из базы данных и писать в топик? Connect. Читать из топика и складывать в S3? Тоже Connect. Писать продюсер или консьюмер вручную для каждой такой задачи — дорого и долго, а Connect делает это декларативно через конфиги.
Запустить Connect можно двумя способами: в standalone-режиме (этот урок) или в distributed-режиме (урок 32). Начинаем с standalone — проще устроен, быстрее поднимается, удобен для знакомства с фреймворком и для разработки.
Что делает connect-standalone.sh
Скрипт connect-standalone.sh запускает один воркер Kafka Connect в standalone-режиме. Всё состояние хранится локально на диске, а не в Kafka-топиках. Коннекторы работают в одном процессе — масштабирования нет, но и сложности настройки нет.
Синтаксис простой: первым аргументом идёт файл конфигурации воркера, следом — один или несколько файлов с конфигурацией коннекторов.
# Проверено: Apache Kafka 4.2.0, Ubuntu 22.04 connect-standalone.sh worker.properties connector1.properties [connector2.properties ...]
Воркер — это сам процесс Connect. Коннекторы — это плагины, которые он загружает и запускает. Один воркер может одновременно запустить несколько коннекторов, если передать несколько файлов.
Архитектура Kafka Connect
Прежде чем настраивать — понять структуру. Connect работает с тремя понятиями: коннектор, таск, воркер.
flowchart LR
subgraph external_src ["Внешняя система (источник)"]
DB["База данных\nфайл / API"]
end
subgraph standalone ["connect-standalone (один процесс)"]
SC["Source Connector\n(плагин)"]
T1["Task 1"]
T2["Task 2"]
SK["Sink Connector\n(плагин)"]
T3["Task 3"]
end
subgraph kafka ["Apache Kafka"]
TP["Топик"]
end
subgraph external_dst ["Внешняя система (приёмник)"]
S3["S3 / HDFS\nElasticsearch / БД"]
end
DB --> SC
SC --> T1
SC --> T2
T1 --> TP
T2 --> TP
TP --> T3
T3 --> SK
SK --> S3
Коннектор — это логическая единица интеграции. Он описывает, откуда читать или куда писать. Таск — единица параллелизма: коннектор делится на таски, каждый таск работает в отдельном потоке. В standalone-режиме всё это живёт в одном JVM-процессе.
Source connector читает из внешней системы и пишет в Kafka. Sink connector читает из Kafka и пишет во внешнюю систему. Можно запустить оба одновременно.
Практическая настройка Kafka connect
Шаг 1. Конфигурация воркера
Файл конфигурации воркера описывает, как Connect подключается к Kafka и где хранит своё состояние. В поставке Kafka есть готовый шаблон — config/connect-standalone.properties. Разберём ключевые параметры.
# config/connect-standalone.properties # Адрес брокеров - обязательно bootstrap.servers=localhost:9092 # Конвертеры для ключей и значений сообщений key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.storage.StringConverter # Для JSON-данных используйте: # 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 # Где хранить offset'ы коннекторов (только для standalone) offset.storage.file.filename=/tmp/connect.offsets # Интервал сброса offset'ов на диск (мс) offset.flush.interval.ms=10000 # REST API воркера - для проверки статуса listeners=HTTP://localhost:8083 # Путь к директории с плагинами коннекторов plugin.path=/usr/share/java,/opt/kafka/plugins
Параметр plugin.path важен: Connect ищет плагины именно здесь. Встроенные коннекторы (FileStream) уже есть в дистрибутиве. Сторонние коннекторы (JDBC, S3, Elasticsearch) нужно скачивать и класть в эту директорию.
В standalone-режиме offset.storage.file.filename указывает локальный файл для хранения offset’ов. Это и есть главное отличие от distributed-режима, где offset’ы хранятся в Kafka-топиках.
Apache Kafka: администрирование кластера
Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Шаг 2. Конфигурация source-коннектора
Для примера используем встроенный FileStreamSource — он читает строки из файла и пишет их в топик. Простейший source-коннектор, не требует сторонних плагинов.
# config/connect-file-source.properties # Уникальное имя коннектора в рамках воркера name=local-file-source # Класс коннектора connector.class=FileStreamSource # Количество тасков tasks.max=1 # Файл-источник file=/tmp/test-source.txt # Топик для записи topic=connect-file-test
Создадим файл-источник с тестовыми данными:
echo "line one" > /tmp/test-source.txt echo "line two" >> /tmp/test-source.txt echo "line three" >> /tmp/test-source.txt
Шаг 3. Конфигурация sink-коннектора
Второй встроенный коннектор — FileStreamSink. Читает из топика и записывает строки в файл. Запустим его вместе с source, чтобы увидеть полный цикл.
# config/connect-file-sink.properties name=local-file-sink connector.class=FileStreamSink tasks.max=1 # Топик для чтения (тот же, что пишет source) topics=connect-file-test # Файл-приёмник file=/tmp/test-sink.txt
Шаг 4. Запуск connect-standalone.sh
Создадим топик заранее — Connect создаст его сам, но лучше контролировать параметры явно:
kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create \ --topic connect-file-test \ --partitions 1 \ --replication-factor 1
Теперь запускаем воркер с обоими коннекторами:
connect-standalone.sh \ config/connect-standalone.properties \ config/connect-file-source.properties \ config/connect-file-sink.properties
В логах должны появиться строки о старте коннекторов и тасков. Процесс не демонизируется — он работает на переднем плане. Для фоновой работы используйте nohup или systemd.
# Запуск в фоне с записью лога nohup connect-standalone.sh \ config/connect-standalone.properties \ config/connect-file-source.properties \ config/connect-file-sink.properties \ > /var/log/kafka/connect-standalone.log 2>&1 & echo "PID: $!"
Шаг 5. Проверка статуса через REST API
Connect предоставляет HTTP-эндпоинты на порту 8083 (или том, что задан в listeners). Через них можно смотреть статус коннекторов без перезапуска процесса.
# Список активных коннекторов curl -s http://localhost:8083/connectors | python3 -m json.tool # Статус конкретного коннектора curl -s http://localhost:8083/connectors/local-file-source/status | python3 -m json.tool # Конфигурация коннектора curl -s http://localhost:8083/connectors/local-file-source/config | python3 -m json.tool # Статус тасков curl -s http://localhost:8083/connectors/local-file-source/tasks | python3 -m json.tool
Ответ на запрос статуса выглядит примерно так:
{
"name": "local-file-source",
"connector": {
"state": "RUNNING",
"worker_id": "localhost:8083"
},
"tasks": [
{
"id": 0,
"state": "RUNNING",
"worker_id": "localhost:8083"
}
],
"type": "source"
}
Если таск завис в состоянии FAILED, там же будет поле trace с трейсом ошибки. Это сильно упрощает отладку.
Шаг 6. Проверка результата
Пока воркер работает, добавим строки в исходный файл и проверим, что они появились в топике и в файле-приёмнике:
# Добавляем данные в источник echo "line four" >> /tmp/test-source.txt echo "line five" >> /tmp/test-source.txt # Читаем из топика напрямую kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic connect-file-test \ --from-beginning \ --max-messages 10 # Проверяем файл-приёмник cat /tmp/test-sink.txt
Задержка между записью в файл и появлением в топике зависит от настройки offset.flush.interval.ms. По умолчанию 10 секунд — это нормально для разработки.
Управление через REST API
В standalone-режиме REST API позволяет не только смотреть статус, но и управлять коннекторами без перезапуска воркера. Это удобно при разработке.
# Остановить таск коннектора
curl -X POST http://localhost:8083/connectors/local-file-source/tasks/0/restart
# Поставить коннектор на паузу
curl -X PUT http://localhost:8083/connectors/local-file-source/pause
# Возобновить
curl -X PUT http://localhost:8083/connectors/local-file-source/resume
# Удалить коннектор (воркер продолжит работу)
curl -X DELETE http://localhost:8083/connectors/local-file-source
# Добавить новый коннектор без перезапуска (через POST)
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
-d '{
"name": "new-source",
"config": {
"connector.class": "FileStreamSource",
"tasks.max": "1",
"file": "/tmp/another.txt",
"topic": "another-topic"
}
}'
Важный момент: коннекторы, добавленные через REST API в standalone-режиме, не сохраняются на диск постоянно. После перезапуска воркера они пропадут — в отличие от distributed-режима, где конфигурация хранится в Kafka.
Когда использовать standalone-режим
Standalone хорошо подходит для нескольких задач. Разработка и тестирование коннекторов — быстрый старт без развёртывания кластера. Одноразовые миграции данных — запустил, перенёс, остановил. Лёгкие сервисы с единственным коннектором, где отказоустойчивость некритична.
В production standalone не используется, потому что нет отказоустойчивости: упал процесс — остановились все коннекторы. Нет масштабирования: добавить воркеров нельзя. Нет хранения конфигурации в Kafka — при рестарте нужно заново передавать файлы свойств. Если нужна продовая интеграция — это distributed. Разбираем его в следующем уроке. Практика настройки production-коннекторов для JDBC, S3, Elasticsearch и мониторинг Connect-кластера входят в программу курса «Apache Kafka для инженеров данных». Там же разбирают написание собственных коннекторов на Java.
Apache Kafka: администрирование кластера
Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Что дальше
В уроке 32 переходим к connect-distributed.sh. Тот же фреймворк, но совсем другая история: несколько воркеров, конфигурация в Kafka, REST API на каждом узле, автоматический ребаланс тасков при падении воркера. Это и есть production-режим Kafka Connect.
Референсы
- 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 (2025)
- KIP-910: Kafka Connect API improvements — cwiki.apache.org (2025)
