Вы сделали грамотную оптимизацию ch, но сам инструмент выбран неверно.
Как бы я это делал во Flink.
Примерно: Kafka exemplar_movements → keyBy(exemplar_id) → поддерживаем enrichment/history state и Kafka exemplar_stocks → keyBy(exemplar_id) → keyed enrichment/join → Kafka enriched_stocks
Главное отличие: не надо пять раз в день загружать несколько миллиардов строк в CH, затем физически перекладывать их ради locality, создавать 5000 jobs, делать JOIN, сохранять промежуточный результат и потом снова читать его из CH в Kafka. Flink сам делает partitioning по ключу и направляет одинаковые exemplar_id в один key-group/operator. Поэтому compute приходит к данным/state, вместо того чтобы заставлять OLAP-БД изображать distributed processing engine. Причём авторы сами очень хорошо показывают симптом: первоначальный CH JOIN занимал ~12 часов, после ручной организации locality они получили ~33 минуты, из которых примерно 13 минут — загрузка, 5 — enrichment и 15 — экспорт. Но есть важная оговорка. У них история — десятки миллиардов записей, плюс нужен расчёт на произвольную прошлую дату. Я бы не пытался тупо засунуть всю эту историю в обычный Flink keyed state. Тут уже интереснее сделать гибрид:
Kafka → Flink → Paimon/Iceberg для полной истории ↘ Flink state для hot/current state stocks → Flink enrichment → Kafka ↘ ClickHouse для аналитики
Вот такая архитектура мне кажется существенно естественнее. И самое забавное: фраза из статьи «50 000 виртуальных партиций дают нам 5 000 независимых задач… упал pod — перезапускаем только его задачи» — это практически описание проблемы, которую Flink уже решает своей моделью key groups + parallelism + checkpointing.
Java укатала всех, включая C. Менее 2 секунд! Понятно, что там глубокая оптимизация. И unsafe, и SIMD, и SWAR, и прочие зверушки... Так что, как серверная экосистема, я считаю, эта одна из лучших.
Очень просто. У Вас есть миллиард объектов. Вам нужно найти объект, удовлетворяющий заданным условиям (состоянию). Прежде чем начать поиск, при помощи фильтра Блума Вы можете определить (с определенной вероятностью), что этот объект есть в этом наборе. Или его там нет. Читал, что фильтр Блума использует Google.
Вы сделали грамотную оптимизацию ch, но сам инструмент выбран неверно.
Как бы я это делал во Flink.
Примерно: Kafka exemplar_movements → keyBy(exemplar_id) → поддерживаем enrichment/history state и Kafka exemplar_stocks → keyBy(exemplar_id) → keyed enrichment/join → Kafka enriched_stocks
Главное отличие: не надо пять раз в день загружать несколько миллиардов строк в CH, затем физически перекладывать их ради locality, создавать 5000 jobs, делать JOIN, сохранять промежуточный результат и потом снова читать его из CH в Kafka. Flink сам делает partitioning по ключу и направляет одинаковые exemplar_id в один key-group/operator. Поэтому compute приходит к данным/state, вместо того чтобы заставлять OLAP-БД изображать distributed processing engine. Причём авторы сами очень хорошо показывают симптом: первоначальный CH JOIN занимал ~12 часов, после ручной организации locality они получили ~33 минуты, из которых примерно 13 минут — загрузка, 5 — enrichment и 15 — экспорт. Но есть важная оговорка. У них история — десятки миллиардов записей, плюс нужен расчёт на произвольную прошлую дату. Я бы не пытался тупо засунуть всю эту историю в обычный Flink keyed state. Тут уже интереснее сделать гибрид:
Kafka → Flink → Paimon/Iceberg для полной истории ↘ Flink state для hot/current state stocks → Flink enrichment → Kafka ↘ ClickHouse для аналитики
Вот такая архитектура мне кажется существенно естественнее. И самое забавное: фраза из статьи «50 000 виртуальных партиций дают нам 5 000 независимых задач… упал pod — перезапускаем только его задачи» — это практически описание проблемы, которую Flink уже решает своей моделью key groups + parallelism + checkpointing.
Тут как-то был челлендж: обработать файл с миллиардом записей. Вот ссылка на ютуб: https://youtu.be/_w4-BqeeC0k?si=73agjjRk5XtEx7HS
Java укатала всех, включая C. Менее 2 секунд! Понятно, что там глубокая оптимизация. И unsafe, и SIMD, и SWAR, и прочие зверушки... Так что, как серверная экосистема, я считаю, эта одна из лучших.