Экосистема Kafka Connect: коннекторы, exactly‑once, DLQ, ретраи

Kafka Connect упрощает интеграцию данных между Kafka и внешними системами без написания кода. Рассказываем, что такое Kafka Connect и как он работает.

Этот текст написала нейросеть, а потом его проверил эксперт и заботливо отредактировал редактор Яндекс Практикума

Что такое Kafka Connect

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

Kafka Connect изучают на курсе «Инженер данных». На нём за 6,5 мес. вы освоите ключевой стек для старта в профессии и научитесь разрабатывать архитектуру данных под бизнес-задачи. После обучения — получите диплом о профессиональной переподготовке.

Зачем Kafka Connect нужен в data-инфраструктуре

В современных data-платформах данные поступают из десятков источников и уходят в десятки систем. Писать код для каждого соединения — дорого и трудно с точки зрения поддержки. Kafka Connect решает эту проблему: он берёт на себя всё управление интеграцией.

Фреймворк автоматически управляет смещениями, обрабатывает ошибки и масштабируется горизонтально. Это избавляет data-команды от рутины и позволяет сосредоточиться на бизнес-логике, а не на написании бесконечных коннекторов.

Как устроена экосистема Kafka Connect

Экосистема Kafka Connect состоит из четырёх основных компонентов. Они работают вместе, чтобы обеспечить надёжную передачу данных. Ниже перечислены ключевые элементы с пояснениями.

  • Воркеры (Workers). Это процессы, которые выполняют подключение и передачу задачи. Воркеры могут работать в двух режимах:
  1. standalone — один процесс, для разработки;
  2. 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 мес. вы освоите ключевой стек для старта в профессии и научитесь разрабатывать архитектуру данных под бизнес-задачи. После обучения — получите диплом о профессиональной переподготовке.

Типовые сценарии применения Kafka Connect

Kafka Connect подходит для многих задач интеграции данных. Он помогает наладить поток данных между системами без написания кода. Вот самые частые сценарии использования:

  • CDC из баз данных. Source-коннекторы захватывают изменения из PostgreSQL, MySQL в реальном времени.
  • Наполнение аналитических хранилищ. Sink-коннекторы отправляют данные из Kafka в Snowflake, BigQuery или ClickHouse.
  • Индексирование данных. Sink-коннекторы пишут в Elasticsearch для обеспечения быстрого поиска.
  • Миграция данных. Source-коннекторы читают из старой системы, sink-коннекторы пишут в новую, и Kafka служит буфером между ними.
  • Интеграция с облачными сервисами. Коннекторы для S3 и Google Cloud Storage для хранения данных в озёрах данных.

Каждый из этих сценариев можно реализовать с помощью готовых коннекторов из каталога. При выборе сценария важно оценить, какие данные критичны для бизнеса и какую задержку допустимо иметь. Например, для CDC критична минимальная задержка, а для миграции — целостность и возможность повтора.

Статью подготовили:
Анатолий Бардуков

ML-инженер в службе качества поиска
Валентина Бокова
Яндекс Практикум
Редактор
Анастасия Павлова
Яндекс Практикум
Иллюстратор

Подпишитесь на наш ежемесячный дайджест статей —
а мы подарим вам полезную книгу про обучение!

Поделиться
Скидка 16% на курсы до 17 сентября. ИИ-навыки уже внутри Забрать скидку
Как ИИ поменяет вашу профессию — пройдите бесплатный тест и получите персональные рекомендации Пройти тест