Приготовьтесь, потому что это будет сложно, но мне необходимо поделиться с вами этим 🙂
🟢 Контекст (дано):
— MongoDB как источник (полуструктурированные данные с гибкой схемой)
— Debezium as Source Connector для репликации лога в Kafka
— MongoDB New Document State Extraction SMT для преобразования сообщений Debezium и загрузки их в DWH
— Конфигурация
array.encoding=document для репликации массивов— Dead Letter Queue (DLQ) для сбора сообщений с ошибками репликации
— Geo Data: коллекция в MongoDB содежащая георгафические зоны в формате GeoJSON
🔴 Проблематика:
— Специфика Array encoding для MongoDB
By default, the event flattening SMT converts MongoDB arrays into arrays that are compatible with Apache Kafka Connect, or Apache Avro schemas. While MongoDB arrays can contain multiple types of elements, all elements in a Kafka array must be of the same type.
— Но без конфигурации
array.encoding=document было бы пропущено по оценкам 50-60% всех реплицируемых данных (недопустимо!)— Чтобы прочесть валидный GeoJSON в СУБД
TRY_TO_GEOGRAPHY(geojson) (в тип GEOGRPAHY) мне нужно реплицировать именно массивы!— Однако и в этом случае теряется несколько записей. Из DLQ я вижу, что это за строки и описание причины:
Field coordinates of schema kafka_mongodbt_connector__geo_zones.geometry.features.geometry is not a homogenous array.
Check option 'struct' of parameter 'array.encoding'
🔵 Решение:
— С помощью Kafka Connect predicates устанавливаем
array.encoding=array для избранных коллекций; Для всех остальных коллекций - array.encoding=document— Анализируем проблемные записи в Dead Letter Queue. Выясняется, что GeoJSON содержит дополнительный ключ
features, в котором дублируются координаты из 3-х чисел, 2 из которых float, а третье всегда 0. И с вероятностью 99% это вызывает ошибку. Необходимо либо проставить 0.0, либо вообще убрать ключ geometry.features из репликации.— Исключаем nested field из репликации на уровне Source Connector:
"field.exclude.list": "db.geo.geometry.features"— Обновляем конфиг работающего коннектора через REST API:
jq .config ${SOURCE__MONGODB} | http PUT ${SOURCE__MONGODB}/config— Инициируем Incremental Snapshot для коллекции с geo через signalling kafka topic
✅ Вуаля! Все записи доступны. Все GeoJson атрибуты корректны.
🟤 Ключевые выводы:
— Допустимо использовать различные варианты конфига
array.encoding (document / array) для разных таблиц (коллекций) БД с помощью Kafka Connect predicates— Без топика Dead Letter Queue (доступен только для Sink Connectors!) не узнал бы об ошибках репликации и их причинах
— Dead Letter Queue удобно просматривать там же в DWH (
SELECT * FROM DLQ), настраивать тесты и уведомления (поэтому для DLQ нужен отдельный Sink Connector!)— Фича Incremental Snapshot помогает прочесть и реплицировать данные по запросу (не останавливая streaming - чтение лога БД / OpLog!)
💬 Задавайте вопросы / оставляйте комментарии.
❤ Кстати, есть тут те, кто тоже хочет real time data streaming? Для целей аналитики, ML, Anti-fraud или прочего оперативного реагирования?
🌐 @data_apps | Навигация по каналу