Kafka Connect на собеседовании Data Engineer

Проверь себя · 1/3разбор после ответа
Нужно получить количество заказов по паре (user_id, status) из таблицы orders. Какой запрос верный?

Что такое Kafka Connect

Kafka Connect — это фреймворк для перекачки данных между Kafka и внешними системами без написания кода. Вместо того чтобы каждый раз руками писать продюсера или консьюмера, вы описываете коннектор в JSON-конфиге и отправляете его в кластер через REST API. Дальше Connect сам вычитывает данные из источника или пишет их в приёмник, следит за смещениями и переживает падения воркеров.

Источник (БД) → Source-коннектор → Kafka → Sink-коннектор → Приёмник (БД)

На собесе Kafka Connect всплывает почти всегда, когда речь заходит про CDC (Change Data Capture) и интеграцию хранилищ: это стандартный способ строить пайплайны «база → Kafka → аналитика» без самописного кода. Интервьюер обычно проверяет, понимаете ли вы разницу между source и sink, зачем нужен distributed mode и как коннектор дробит работу на задачи (tasks).

Source-коннекторы

Source-коннектор читает данные из внешней системы и публикует их в топики Kafka.

  • Debezium — самый популярный source-коннектор для CDC. Читает журнал изменений базы (WAL в Postgres, binlog в MySQL) и превращает каждую вставку, обновление и удаление в событие. Поддерживает Postgres, MySQL, MongoDB, Oracle, SQL Server. Главное преимущество — ловит изменения почти в реальном времени и не нагружает базу лишними запросами.
  • JDBC source — периодически выполняет SELECT из реляционной базы по расписанию. Проще Debezium, но менее эффективен: он опрашивает таблицу, а не читает журнал, поэтому создаёт нагрузку запросами и плохо ловит удаления и обновления (нужен инкрементный столбец или timestamp).
  • File source — следит за директорией и публикует содержимое файлов в топик. Обычно используется для простых сценариев и демо.
  • S3 source — читает файлы из объектного хранилища S3 и заливает их в Kafka.
{
  "name": "postgres-cdc",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "...",
    "database.dbname": "...",
    "schema.include": "public",
    "table.include.list": "public.orders,public.users"
  }
}

Sink-коннекторы

Sink-коннектор делает обратное: вычитывает сообщения из топиков Kafka и пишет их во внешнюю систему.

  • Elasticsearch sink — индексирует данные из топика для полнотекстового поиска.
  • JDBC sink — пишет сообщения обратно в реляционную базу. Для защиты от дублей обычно настраивают идемпотентную запись (UPSERT по ключу).
  • S3 sink — архивирует поток в S3 в формате Parquet или JSON, обычно с ротацией файлов по времени или размеру. Частый способ дёшево складывать сырые события в data lake.
  • ClickHouse sink — стримит события напрямую в ClickHouse для аналитики.
  • HDFS sink — выгружает данные в HDFS для батч-обработки.
{
  "name": "s3-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "topics": "orders",
    "s3.bucket.name": "my-bucket",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "rotate.schedule.interval.ms": "60000"
  }
}

Distributed mode

У Kafka Connect два режима запуска, и разницу между ними любят спрашивать на собесе.

  • Standalone. Один воркер, вся конфигурация в файле. Подходит для разработки, тестов и небольших локальных задач. Нет отказоустойчивости: упал процесс — встал пайплайн.
  • Distributed. Несколько воркеров образуют кластер и делят работу между собой. Конфигурация, смещения и статусы коннекторов хранятся в служебных топиках Kafka, а не в файле, поэтому кластер переживает падение отдельного воркера. Управляют коннекторами через REST API, а масштабируется всё горизонтально — добавлением воркеров.

Коннектор дробит работу на задачи (tasks), а tasks.max ограничивает их число. Именно задачи распределяются по воркерам и дают параллелизм: например, коннектор на десять таблиц может разложить их на несколько задач.

Воркер 1: Task1, Task2 (коннектор A)
Воркер 2: Task3 (коннектор B), Task4 (коннектор A)

Когда воркер падает, кластер запускает ребаланс и перераспределяет его задачи на живые воркеры — пайплайн продолжает работать.

Готовься к собесу аналитика как в Duolingo
10 минут в день — SQL, Python, A/B, метрики. 1700+ вопросов в Telegram
Открыть Карьерник в Telegram

Single Message Transforms

Single Message Transforms (SMT) — лёгкие преобразования отдельных сообщений прямо внутри Connect, без написания кода. Работают строго над одним сообщением: агрегации, джойны и обращения к внешним системам им недоступны — для этого нужен полноценный стрим-процессинг (Kafka Streams, Flink).

  • Маскирование полей — скрыть чувствительные данные (номер карты, СНИЛС) перед публикацией.
  • Переименование полей — привести схему к нужному виду.
  • Приведение типов — например, строку в число.
  • Фильтрация сообщений — отбросить события, не подходящие под условие.
  • Добавление метки времени — проставить служебное поле.
"transforms": "MaskField",
"transforms.MaskField.type": "org.apache.kafka.connect.transforms.MaskField$Value",
"transforms.MaskField.fields": "ssn,credit_card"

Не путайте SMT с конвертерами (key.converter / value.converter): конвертеры отвечают за формат сериализации сообщения (Avro, JSON, Protobuf), а SMT — за преобразование его содержимого. На собесе это частый уточняющий вопрос.

Частые ошибки

  • Путать source и sink. Source читает из внешней системы в Kafka, sink — из Kafka во внешнюю систему. Ошибиться на этом — сразу минус.
  • Считать JDBC source полноценным CDC. JDBC опрашивает таблицу запросами и плохо ловит удаления; настоящий CDC — это Debezium, который читает журнал базы.
  • Ждать от SMT сложной логики. SMT — это преобразования одного сообщения, без агрегаций и джойнов. Если нужно объединять потоки, это уже Kafka Streams или Flink, а не Connect.
  • Использовать standalone в проде. Без distributed mode нет отказоустойчивости: падение воркера останавливает пайплайн.
  • Игнорировать обработку ошибок. Без настройки errors.tolerance и dead letter queue одно битое сообщение способно уронить задачу коннектора.

Связанные темы

FAQ

Чем Kafka Connect отличается от самописного продюсера или консьюмера?

Connect — это готовый фреймворк: вы описываете интеграцию конфигом, а он берёт на себя чтение смещений, распределение задач, отказоустойчивость и рестарты. Самописный код придётся поддерживать самому, зато он даёт полный контроль. Connect выигрывает, когда нужна типовая интеграция «база ↔ Kafka» без изобретения велосипеда.

В чём разница между source- и sink-коннектором?

Source-коннектор читает данные из внешней системы и публикует их в Kafka (например, Debezium из Postgres). Sink-коннектор вычитывает из топиков Kafka и пишет во внешнюю систему (например, в S3 или ClickHouse).

Чем Debezium лучше JDBC source для CDC?

Debezium читает журнал изменений базы (WAL/binlog) и ловит вставки, обновления и удаления почти в реальном времени, не нагружая базу запросами. JDBC source периодически опрашивает таблицу через SELECT, создаёт нагрузку и плохо отслеживает удаления. Для настоящего CDC используют Debezium.

Что можно и чего нельзя делать через SMT?

Через SMT можно преобразовывать отдельные сообщения: маскировать и переименовывать поля, приводить типы, фильтровать, добавлять метки времени. Нельзя делать агрегации, джойны потоков и обращаться к внешним системам — для этого нужен Kafka Streams или Flink.

Зачем нужен distributed mode?

Он даёт отказоустойчивость и горизонтальное масштабирование: воркеры образуют кластер, хранят конфигурацию и смещения в служебных топиках Kafka и перераспределяют задачи упавшего воркера на живые. Standalone-режим этого не умеет и годится только для разработки.

Это официальная информация?

Нет. Статья основана на документации Apache Kafka, Confluent и Debezium и на опыте кандидатов. Конкретные требования зависят от компании, команды и уровня позиции.


Тренируйте Data Engineering — откройте тренажёр с 1500+ вопросами для собесов.