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

Изучаем Apache Kafka с нуля. Урок 31. connect-standalone.sh: Kafka Connect в 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.

Референсы

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

Тема 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/