Зачем Kafka Connect нужен в data-инфраструктуре
В современных data-платформах данные поступают из десятков источников и уходят в десятки систем. Писать код для каждого соединения — дорого и трудно с точки зрения поддержки. Kafka Connect решает эту проблему: он берёт на себя всё управление интеграцией.
Фреймворк автоматически управляет смещениями, обрабатывает ошибки и масштабируется горизонтально. Это избавляет data-команды от рутины и позволяет сосредоточиться на бизнес-логике, а не на написании бесконечных коннекторов.
Как устроена экосистема Kafka Connect
Экосистема Kafka Connect состоит из четырёх основных компонентов. Они работают вместе, чтобы обеспечить надёжную передачу данных. Ниже перечислены ключевые элементы с пояснениями.
- Воркеры (Workers). Это процессы, которые выполняют подключение и передачу задачи. Воркеры могут работать в двух режимах:
- standalone — один процесс, для разработки;
- distributed — кластер из нескольких процессов, для продакшена, отказоустойчивый и масштабируемый.
- Коннекторы (Connectors). Это конфигурации, которые определяют, как подключаться к внешней системе и какие данные передавать.
- Задачи (Tasks). Это реальные исполнители, которые копируют данные. Коннектор разбивает работу на несколько задач, которые выполняются параллельно на разных воркерах.
- Конвертеры (Converters). Они отвечают за сериализацию и десериализацию данных. Конвертеры «переводят» данные из внутреннего формата Kafka Connect в байтовый массив (который фактически и хранится в топике Kafka) и обратно.
Понимание этих компонентов помогает правильно настраивать и отлаживать коннекторы. Каждый элемент можно конфигурировать отдельно — это даёт гибкость для разных сценариев. Например, для работы с Avro-схемами нужен правильный конвертер, а для высокой нагрузки — достаточное количество воркеров.
Source- и sink-коннекторы: в чём разница
Коннекторы делятся на source- и sink-коннекторы. Source-коннекторы забирают данные из внешних систем и публикуют их в топики Kafka — это «входные ворота» в экосистему. Sink-коннекторы читают данные из топиков Kafka и отправляют их во внешние системы — это «выходные ворота» из экосистемы. Ниже в таблице показаны ключевые различия.
| Характеристики |
Source-коннектор |
Sink-коннектор |
| Направление |
Внешняя система → Kafka |
Kafka → внешняя система |
| Роль в Kafka |
Продюсер |
Консьюмер |
| Управление смещениями |
Хранит свой прогресс в offset.storage.topic |
Использует стандартный механизм consumer group |
| Параллелизм |
Зависит от источника. Например, несколько файлов или партиций БД |
Ограничен числом партиций в топиках, которые читает коннектор |
Как работают коннекторы Kafka Connect
Source-коннекторы регулярно опрашивают источник (polling) или слушают изменения (CDC). Они сохраняют информацию о том, какие данные уже отправлены, чтобы не дублировать записи при перезапуске. Sink-коннекторы работают как обычные консьюмеры: они читают записи из топиков и отправляют их в целевую систему, используя стандартные смещения consumer group.
Для sink-коннекторов количество задач ограничено числом партиций в топике: одна партиция не может обрабатываться более чем одной задачей одновременно. Если коннектор читает несколько топиков, максимальное число активных задач определяется по топику с наибольшим количеством партиций.
Exactly-once semantics: что это значит и где работает
Exactly-once — это гарантия того, что каждая запись будет доставлена ровно один раз даже при сбоях. Это идеальный сценарий, но он зависит от возможностей конкретного коннектора. Для source-коннекторов exactly-once настраивается через свойство exactly.once.source.support на уровне воркера.
На практике не все коннекторы поддерживают exactly-once. Некоторые коннекторы могут гарантировать только at-least-once, но с идемпотентными записями это даёт тот же эффект. Разработчики намеренно избегают термина «exactly-once delivery» в коде и документации, подчёркивая, что это зависит не только от фреймворка, но и от конкретного коннектора.
DLQ: зачем нужна dead letter queue
Dead Letter Queue (DLQ) — это отдельный топик, куда отправляются записи, которые не удалось обработать. Это позволяет сохранить проблемные записи для последующего анализа — вместо того чтобы потерять их или бесконечно пытаться обработать. DLQ работает для ошибок на этапе конвертации и трансформации данных.
В sink-коннекторах можно настроить политику обработки ошибок: NOOP (игнорировать), THROW (остановить коннектор) или RETRY (повторять с настраиваемым интервалом). Важное ограничение: DLQ не применяется к ошибкам внутри метода put() sink-коннектора — при ошибках записи в целевую систему записи не попадают в DLQ.
Ретраи и обработка ошибок в Kafka Connect
Ретраи — это повторные попытки выполнить операцию после ошибки. В sink-коннекторах можно настроить политику RETRY с указанием максимального числа попыток и интервала между ними. Это помогает справляться с временными сбоями сети или перегрузками целевой системы.
Существуют RetriableException и FatalException. RetriableException — это временные ошибки, которые могут быть исправлены повторной попыткой. FatalException — это ошибки, которые не исправляются повторными попытками. Для RetriableException коннектор продолжает ретраить, для FatalException запись отправляется в DLQ.
Kafka Connect изучают на курсе «Инженер данных». На нём за 6,5 мес. вы освоите ключевой стек для старта в профессии и научитесь разрабатывать архитектуру данных под бизнес-задачи. После обучения — получите диплом о профессиональной переподготовке.