Source Connectors
Source connector читает данные из внешней системы и публикует их в Kafka. Без source connector вы бы писали продюсер вручную: подключались к базе данных, запрашивали новые строки, отслеживали позицию, публиковали в Kafka, обрабатывали ошибки. Source connector делает всё это в рамках стандартного фреймворка — с автоматическим управлением offset, перераспределением задач и мониторингом через REST API.
Жизненный цикл source connector
Понимание жизненного цикла — ключ к отладке и правильной конфигурации.
- Инициализация. Worker вызывает
SourceConnector.start(props)— коннектор проверяет конфигурацию и устанавливает соединение. - Разбиение на задачи. Worker вызывает
SourceConnector.taskConfigs(maxTasks)— коннектор возвращает список конфигураций задач (не болееtasks.max). - Распределение задач. Worker распределяет задачи по воркерам кластера.
- Чтение данных. Каждая задача вызывает
SourceTask.poll()в цикле.poll()возвращаетList<SourceRecord>— записи для публикации в Kafka. - Сериализация и публикация. Worker сериализует записи через converter и публикует в Kafka через internal producer.
- Фиксация offset. После успешной публикации Worker записывает offset в топик
connect-offsets. Если Worker падает до фиксации — при перезапуске задача начнёт с последнего зафиксированного offset, возможны дубликаты (at-least-once).
Источник
Внешний источник данных: база данных, файловая система, REST API, очередь сообщений. Source connector абстрагирует специфику конкретного источника за стандартным интерфейсом SourceTask.poll().Source Task
SourceTask: единица работы. Вызывает poll() в цикле, возвращает List<SourceRecord>. Каждый SourceRecord содержит: sourcePartition (откуда), sourceOffset (позиция), topic (куда писать), key, value.Worker
Connect Worker: принимает SourceRecord от Task, сериализует через converter (JSON, Avro, Protobuf), публикует в Kafka через internal producer. После подтверждения от Kafka — записывает offset в connect-offsets.Kafka Topic
Kafka Topic: целевой топик, указанный в конфигурации коннектора. Prefix topic.prefix + table name для JDBC, или явно заданный topic для FileStream.FileStreamSourceConnector: учебный пример
FileStreamSourceConnector — простейший source connector. Читает файл построчно, каждая строка = одна запись в Kafka.
{
"name": "file-source-example",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"file": "/var/log/app.log",
"topic": "app-logs",
"tasks.max": "1"
}
}
Offset для FileStream — это позиция байта в файле. При перезапуске коннектор продолжит чтение с последней зафиксированной позиции.
FileStreamSourceConnector предназначен исключительно для учебных примеров и тестирования. В production не использовать: нет мониторинга ротации логов, нет поддержки сжатых файлов, нет отказоустойчивости для файловых систем. Для production-сценариев работы с файлами используйте Kafka Connect S3 Source или Splunk Kafka Connector.
JDBC Source Connector: чтение из реляционных баз
JdbcSourceConnector читает данные из реляционных баз данных через JDBC. Поддерживает PostgreSQL, MySQL, Oracle, SQL Server, и любую базу с JDBC-драйвером.
Конфигурация:
{
"name": "jdbc-orders-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:postgresql://db:5432/shop",
"connection.user": "kafka_user",
"connection.password": "secret",
"table.whitelist": "orders,order_items",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "db-",
"poll.interval.ms": "5000",
"tasks.max": "2"
}
}
Это создаст топики db-orders и db-order_items.
Режимы чтения JDBC Source
Параметр mode определяет, как коннектор отслеживает уже прочитанные строки:
| Режим | Механизм отслеживания | Применение |
|---|---|---|
bulk | Нет отслеживания — полная выгрузка каждый poll | Небольшие справочные таблицы, которые меняются редко |
incrementing | Монотонно возрастающий числовой ID (incrementing.column.name) | Таблицы с автоинкрементным первичным ключом, только INSERT |
timestamp | Столбец с датой последнего изменения (timestamp.column.name) | Таблицы с updated_at, поддерживают UPDATE |
timestamp+incrementing | Комбинация timestamp и ID | Таблицы с UPDATE и INSERT, наиболее точный режим |
Режим incrementing запоминает максимальный ID, прочитанный на предыдущем poll. Следующий запрос: SELECT * FROM orders WHERE id > :last_id ORDER BY id. Это надёжно для INSERT, но не видит UPDATE и DELETE.
Режим timestamp запоминает максимальное значение updated_at. Следующий запрос: SELECT * FROM orders WHERE updated_at > :last_timestamp. Видит UPDATE, но может пропустить строки при одинаковом timestamp с точностью до секунды.
JDBC Source не поддерживает DELETE. Если строка удалена из базы — коннектор это не видит. Для полного CDC (включая DELETE) используйте Debezium.
Debezium: Change Data Capture как Source Connector
Debezium — это набор Kafka Connect source connectors для захвата изменений базы данных (CDC — Change Data Capture). В отличие от JDBC Source, который периодически опрашивает таблицы, Debezium читает бинарный лог базы данных (binlog MySQL, WAL PostgreSQL, oplog MongoDB) и фиксирует каждое изменение в реальном времени.
Поддерживаемые базы данных: PostgreSQL, MySQL, MongoDB, SQL Server, Oracle, Db2, Cassandra.
Минимальная конфигурация Debezium PostgreSQL connector:
{
"name": "debezium-postgres-source",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "debezium",
"database.password": "dbz",
"database.dbname": "mydb",
"database.server.name": "mydb-server",
"table.include.list": "public.orders,public.customers",
"plugin.name": "pgoutput",
"tasks.max": "1"
}
}
Debezium создаёт топики с именами {server.name}.{schema}.{table}. Для конфигурации выше: mydb-server.public.orders, mydb-server.public.customers.
Структура события Debezium
Каждое событие Debezium содержит не просто новую строку, а полный контекст изменения:
{
"before": { "id": 1, "status": "pending" },
"after": { "id": 1, "status": "shipped" },
"op": "u",
"ts_ms": 1715000000000,
"source": {
"db": "mydb",
"table": "orders",
"lsn": 12345678
}
}
Поле op: c (create/INSERT), u (update/UPDATE), d (delete/DELETE), r (read/snapshot).
Debezium vs JDBC Source
| Характеристика | JDBC Source | Debezium CDC |
|---|---|---|
| Задержка | Минуты (poll interval) | Миллисекунды (log streaming) |
| DELETE | Не видит | Видит (op: "d") |
| UPDATE | Только с timestamp mode | Всегда (before + after) |
| Нагрузка на БД | Периодические SELECT-запросы | Минимальная (чтение WAL/binlog) |
| Snapshot | Нет | Начальный снимок при первом запуске |
Debezium — это полноценный курс на нашей платформе. Здесь мы рассматриваем его как Source Connector в контексте Kafka Connect. Для глубокого изучения CDC, репликации WAL, трансформации событий и Debezium в production — обратитесь к отдельному курсу.
Offset tracking: как source connector запоминает позицию
Source connector сохраняет свою позицию в Kafka-топике connect-offsets. Механизм зависит от типа коннектора:
| Коннектор | Offset (что хранится) | Пример значения |
|---|---|---|
| FileStream | Байтовая позиция в файле | {"position": 4096} |
| JDBC incrementing | Максимальный ID | {"incrementing": 1042} |
| JDBC timestamp | Максимальный timestamp | {"timestamp": 1715000000000} |
| Debezium PostgreSQL | LSN (Log Sequence Number) WAL | {"lsn": 12345678} |
| Debezium MySQL | Binlog filename + position | {"file": "mysql-bin.000001", "pos": 4096} |
Ключ в connect-offsets состоит из имени коннектора и source partition. Значение — JSON с текущей позицией.
Если connect-offsets топик по каким-то причинам утерян (удалён) — все source коннекторы потеряют свою позицию и начнут с начала (или с earliest). Это приведёт к дублированию данных в Kafka. Поэтому connect-offsets должен быть защищён от случайного удаления.
Параллелизм source connectors
Параметр tasks.max определяет максимальное число задач. Реальный параллелизм зависит от поддержки источника:
- FileStream: поддерживает только 1 задачу (файл линейный).
- JDBC Source: может создавать несколько задач при наличии нескольких таблиц в
table.whitelist. Каждая таблица — одна задача. - Debezium: обычно 1 задача на один сервер (1 WAL-поток). Параллелизм достигается через несколько коннекторов для разных баз.
{
"table.whitelist": "orders,customers,products",
"tasks.max": "3"
}
При tasks.max=3 и трёх таблицах JDBC Source создаст три задачи: каждая читает свою таблицу. Задачи распределяются по воркерам Connect-кластера.
Avro converter с Schema Registry
Для source connector, пишущего Avro-данные, необходимо указать AvroConverter и URL Schema Registry:
{
"value.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081"
}
При первой записи converter регистрирует схему в Schema Registry и получает числовой schema ID. Каждое последующее сообщение содержит 5-байтовый заголовок (0x00 + 4-байтовый schema ID) перед Avro-payload. Детально этот механизм рассматривается в Модуле 06 (Schema Registry).
Ключевые выводы
- Source connector читает данные из внешней системы через стандартный интерфейс
SourceTask.poll(). - FileStreamSource — только для учебных примеров.
- JDBC Source поддерживает несколько режимов отслеживания (bulk, incrementing, timestamp). Не видит DELETE.
- Debezium читает бинарный лог базы данных: минимальная задержка, полное CDC (INSERT + UPDATE + DELETE).
- Offset хранится в Kafka-топике
connect-offsets(не в__consumer_offsets— это только для sink). - AvroConverter требует Schema Registry — детали в Модуле 06.