Небольшой учебный проект, который демонстрирует потоковую обработку рекламных событий в «псевдо‑проде»:
- генерирует события (impression, click, conversion);
- сбрасывает их в прочный лог (эмуляция Kafka) с возможностью replay;
- агрегирует метрики по ключу (campaign_id, ad_id) и периодически делает снапшот как бы ClickHouse в TSV;
- отдает runtime‑метрики в формате Prometheus по HTTP (/metrics).
Проект сфокусирован на практических аспектах: консистентность, восстановление после падения, работа с файлами как с лог‑хранилищем, backfill/replay и экспонирование метрик.
-
Producer (внутри main):
- синтетически генерирует события с заданной скоростью EPS;
- пишет в
KafkaTopic(append‑only файл TSV).
-
KafkaTopic/KafkaConsumer (эмуляция Kafka):
KafkaTopic(data/kafka/topic.tsv) — один топик/раздел, только append;KafkaConsumer— tail чтение файла с указанного offset (количество уже записанных строк).
-
ClickHouseSim:
- raw строки складываются в
data/ch/raw/events.tsv; - периодический снапшот агрегатов — атомарная замена файла
data/ch/agg/snapshot.tsvчерез временный файл.
- raw строки складываются в
-
MetricsAggregator:
- потребляет события из лога, считает Counters: impressions, clicks, conversions, cost, revenue;
- производные метрики: CTR, CPC, ROI;
- каждые 1с пишет снапшот и чекпоинт (offset) в
data/checkpoint.txt— это и есть простой механизм replay.
-
PrometheusServer (Winsock):
- поднимает HTTP на
127.0.0.1:<port>и отдает/metricsв формате Prometheus exposition.
- поднимает HTTP на
Требования: CMake 3.16+, компилятор с C++17.
cmake -S . -B build -DCMAKE_BUILD_TYPE=Release
cmake --build build --config ReleaseCLion: просто «Build» — CMakeLists уже подключает исходники и Ws2_32.
build\ad_pipeline.exe --duration-seconds 15 --eps 200 --prom-port 9090 --prom-addr 127.0.0.1 --data-dir dataПараметры:
--duration-seconds— длительность демо‑прогона (по умолчанию 15);--eps— events per second (по умолчанию 200);--prom-addrи--prom-port— адрес и порт HTTP /metrics (по умолчанию 127.0.0.1:9090);--data-dir— базовая папка данных (по умолчаниюdata).
За время работы:
- raw log:
data/kafka/topic.tsv; - raw как в ClickHouse:
data/ch/raw/events.tsv; - агрегаты:
data/ch/agg/snapshot.tsv(перезаписывается раз в 1с); - чекпоинты:
data/checkpoint.txt(последний прочитанный offset).
Во время прогона откройте:
http://127.0.0.1:9090/metrics
Примеры метрик:
ad_pipeline_processed_events_total 12345
ad_pipeline_groups 9
ad_pipeline_impressions_total{campaign_id="cmp-1",ad_id="ad-2"} 321
ad_pipeline_clicks_total{campaign_id="cmp-1",ad_id="ad-2"} 12
ad_pipeline_conversions_total{campaign_id="cmp-1",ad_id="ad-2"} 1
ad_pipeline_cost_total{campaign_id="cmp-1",ad_id="ad-2"} 3.450000
ad_pipeline_revenue_total{campaign_id="cmp-1",ad_id="ad-2"} 1.500000
ad_pipeline_ctr{campaign_id="cmp-1",ad_id="ad-2"} 0.037383
ad_pipeline_cpc{campaign_id="cmp-1",ad_id="ad-2"} 0.287500
ad_pipeline_roi{campaign_id="cmp-1",ad_id="ad-2"} -0.565217
Фрагмент prometheus.yml:
scrape_configs:
- job_name: 'ad-pipeline'
scrape_interval: 2s
static_configs:
- targets: ['127.0.0.1:9090']После запуска проекта метрики будут подтягиваться автоматически.
data/ch/agg/snapshot.tsv— актуальное состояние агрегатов (с CTR/CPC/ROI);data/ch/raw/events.tsv— «сырые» события, как бы табличка ClickHouse (TSV);data/kafka/topic.tsv— «лог» событий, откуда идет replay;data/checkpoint.txt— offset, с которого агрегатор продолжит при рестарте.
- Семантика доставки — at‑least‑once; возможен двойной учет при рестарте.
- Один топик/раздел, нет реального партиционирования, нет backpressure.
- HTTP‑сервер минималистичен (один endpoint, примитивный парсер запроса).
- «ClickHouse» — это просто TSV‑файлы; никаких индексов/компрессии/SQL.