Kafka Connect на собеседовании Data Engineer
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)Когда воркер падает, кластер запускает ребаланс и перераспределяет его задачи на живые воркеры — пайплайн продолжает работать.
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 одно битое сообщение способно уронить задачу коннектора.
Связанные темы
- Kafka на собесе DE
- CDC и Debezium для DE
- Schema evolution для DE
- Kafka consumer groups для DE
- Подготовка к собесу Data Engineer
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+ вопросами для собесов.