Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 06.02 · 25 мин
Средний
Source ConnectorFileStreamSourceJdbcSourceDebeziumCDCOffset Tracking

Source Connectors

Source connector читает данные из внешней системы и публикует их в Kafka. Без source connector вы бы писали продюсер вручную: подключались к базе данных, запрашивали новые строки, отслеживали позицию, публиковали в Kafka, обрабатывали ошибки. Source connector делает всё это в рамках стандартного фреймворка — с автоматическим управлением offset, перераспределением задач и мониторингом через REST API.


Жизненный цикл source connector

Понимание жизненного цикла — ключ к отладке и правильной конфигурации.

  1. Инициализация. Worker вызывает SourceConnector.start(props) — коннектор проверяет конфигурацию и устанавливает соединение.
  2. Разбиение на задачи. Worker вызывает SourceConnector.taskConfigs(maxTasks) — коннектор возвращает список конфигураций задач (не более tasks.max).
  3. Распределение задач. Worker распределяет задачи по воркерам кластера.
  4. Чтение данных. Каждая задача вызывает SourceTask.poll() в цикле. poll() возвращает List<SourceRecord> — записи для публикации в Kafka.
  5. Сериализация и публикация. Worker сериализует записи через converter и публикует в Kafka через internal producer.
  6. Фиксация offset. После успешной публикации Worker записывает offset в топик connect-offsets. Если Worker падает до фиксации — при перезапуске задача начнёт с последнего зафиксированного offset, возможны дубликаты (at-least-once).
Жизненный цикл Source Connector

Источник

Внешний источник данных: база данных, файловая система, REST API, очередь сообщений. Source connector абстрагирует специфику конкретного источника за стандартным интерфейсом SourceTask.poll().
poll()

Source Task

SourceTask: единица работы. Вызывает poll() в цикле, возвращает List<SourceRecord>. Каждый SourceRecord содержит: sourcePartition (откуда), sourceOffset (позиция), topic (куда писать), key, value.
SourceRecord

Worker

Connect Worker: принимает SourceRecord от Task, сериализует через converter (JSON, Avro, Protobuf), публикует в Kafka через internal producer. После подтверждения от Kafka — записывает offset в connect-offsets.
produce

Kafka Topic

Kafka Topic: целевой топик, указанный в конфигурации коннектора. Prefix topic.prefix + table name для JDBC, или явно заданный topic для FileStream.
commit offset
connect-offsetsCompacted Kafka-топик для хранения позиций source-коннекторов. Ключ: (connector name, source partition). Значение: текущий offset. При перезапуске задача читает этот топик и продолжает с последней зафиксированной позиции.

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 — это позиция байта в файле. При перезапуске коннектор продолжит чтение с последней зафиксированной позиции.

WARNING

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 с точностью до секунды.

NOTE

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 SourceDebezium CDC
ЗадержкаМинуты (poll interval)Миллисекунды (log streaming)
DELETEНе видитВидит (op: "d")
UPDATEТолько с timestamp modeВсегда (before + after)
Нагрузка на БДПериодические SELECT-запросыМинимальная (чтение WAL/binlog)
SnapshotНетНачальный снимок при первом запуске
NOTE

Debezium — это полноценный курс на нашей платформе. Здесь мы рассматриваем его как Source Connector в контексте Kafka Connect. Для глубокого изучения CDC, репликации WAL, трансформации событий и Debezium в production — обратитесь к отдельному курсу.

Архитектура Debezium: как устроен CDC-коннектор

Offset tracking: как source connector запоминает позицию

Source connector сохраняет свою позицию в Kafka-топике connect-offsets. Механизм зависит от типа коннектора:

КоннекторOffset (что хранится)Пример значения
FileStreamБайтовая позиция в файле{"position": 4096}
JDBC incrementingМаксимальный ID{"incrementing": 1042}
JDBC timestampМаксимальный timestamp{"timestamp": 1715000000000}
Debezium PostgreSQLLSN (Log Sequence Number) WAL{"lsn": 12345678}
Debezium MySQLBinlog 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).


Ключевые выводы

  1. Source connector читает данные из внешней системы через стандартный интерфейс SourceTask.poll().
  2. FileStreamSource — только для учебных примеров.
  3. JDBC Source поддерживает несколько режимов отслеживания (bulk, incrementing, timestamp). Не видит DELETE.
  4. Debezium читает бинарный лог базы данных: минимальная задержка, полное CDC (INSERT + UPDATE + DELETE).
  5. Offset хранится в Kafka-топике connect-offsets (не в __consumer_offsets — это только для sink).
  6. AvroConverter требует Schema Registry — детали в Модуле 06.
Проверка знанийKnowledge check
Таблица `orders` в PostgreSQL имеет столбцы: id (SERIAL), amount (DECIMAL), updated_at (TIMESTAMP). Команда хочет читать новые строки и обновления существующих через JDBC Source Connector. Какой mode выбрать и почему режим incrementing недостаточен?
ОтветAnswer
Нужен режим timestamp+incrementing (или только timestamp). Режим incrementing запоминает максимальный ID и запрашивает строки с id > last_id. Это работает только для новых INSERT: если существующая строка обновляется (UPDATE), её ID не меняется, и коннектор не увидит изменение. Режим timestamp использует updated_at для обнаружения изменений: строки с updated_at > last_timestamp попадают в следующий poll — это захватывает и INSERT, и UPDATE. Режим timestamp+incrementing комбинирует оба подхода: надёжно при одинаковых timestamp у нескольких строк. Для полного CDC (включая DELETE) необходим Debezium.

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 7