Спроектируй агрегацию кликов по рекламе
Система, считающая клики по рекламе для биллинга: потоковая обработка Kafka и Flink, exactly-once через идемпотентность, оконная агрегация, поздние события и watermark, дедупликация и пакетный слой сверки, корректирующий быстрый путь.
Рекламодатель открыл спор: дашборд платформы говорил, что его кампания получила 1,2 миллиона кликов за месяц, а его собственная аналитика лендинга насчитала 1,0 миллиона. Двести тысяч кликов — по несколько центов каждый — выставлялись в счёт ни за что. Команда инженеров покопалась и нашла три причины, сплетённые вместе: сетевой ретрай, отправивший часть событий клика дважды; потоковый процессор, посчитавший дважды после крэша-и-реплея; и пачку кликов, пришедших с опозданием на часы от пользователей на дрянной мобильной связи и выброшенных из окна, к которому они принадлежали. Ни одна из них не экзотика. Это поведение по умолчанию наивного потокового конвейера, и именно поэтому агрегация кликов по рекламе — подсчёт событий за деньги — одна из труднейших задач «просто посчитай штуки» в поле.
Требования
Система принимает поток событий кликов и выдаёт агрегированные счётчики — кликов на рекламу, в минуту, по региону — что питают и real-time дашборды, и систему биллинга. Поскольку выход — деньги, планка куда выше, чем для аналитики.
Функциональные. Принимать события кликов (ad_id, user_id, timestamp, click_id, регион). Агрегировать в временные окна (счётчики на минуту, на час). Обслуживать быстрые запросы («клики за последние 5 минут») для дашборда рекламодателя и выдавать авторитетный дневной итог для биллинга. Фильтровать фрод и невалидные клики. Давать рекламодателям резать по измерениям (реклама, кампания, гео).
Нефункциональные. Корректность первостепенна — каждый биллингуемый клик посчитан ровно один раз, не больше и не меньше, потому что пересчёт переплачивает (споры, возвраты, юридический риск), а недосчёт теряет выручку. Высокая пропускная способность: крупная сеть видит сотни тысяч кликов в секунду на пике. Низкая задержка для дашборда (свежесть в секунды), но биллинг терпит часы задержки в обмен на правильность. Система должна переживать отказы, не портя счётчики.
Напряжение, формирующее всё: клик — это событие во времени, сеть доставляет события не по порядку и иногда дважды, а процессоры падают на полпути. Получить точный счёт из этого месива — вся работа.
Оценка
пик кликов = 1 000 000 кликов/с (крупная сеть на пике)
размер события = ~100 байт на событие клика
пик приёма = 10^6 × 100 байт = 100 МБ/с в лог
кликов в день = ~1 000 000 × 86 400 (средн. ~1/3 пика) ≈ 3 × 10^10/день
сырое хранение = 3 × 10^10 × 100 Б ≈ 3 ТБ/день сырых событий
агрегированный выход = счётчики × окна ≪ сырого (счётчики крошечны против событий)Дизайн ведут два наблюдения. Первое: объём сырых событий огромен (100 МБ/с, много ТБ/день), но агрегированный выход крошечен — работа системы в том, чтобы схлопнуть поток в маленькие счётчики, так что дорогая часть — потоковое вычисление, а не хранение результата. Второе: надо хранить сырые события (в логе / холодном хранилище) даже после агрегации, потому что эта сырая запись — то, что пакетный слой сверки перечитывает, чтобы поправить быстрый путь, и то, что показываешь рекламодателю в споре. Агрегаты производны; события — источник истины.
Высокоуровневый дизайн
Это lambda-архитектура в классической форме: слой скорости (потоковый) для низколатентных-но-приближённых счётчиков дашборда и пакетный слой, перечитывающий неизменяемые сырые события, чтобы посчитать правильные итоги и свериться со слоем скорости. Durable, партиционированный лог (Kafka) — стержень: это воспроизводимый источник истины, который потребляют оба слоя, и воспроизводимость — то, что позволяет пакетному слою пересчитать точно, а потоковому восстановиться после крэша. Современная альтернатива, kappa, выбрасывает отдельный пакетный код и гоняет всё через потоковый движок, перепроцессируя реплеем лога, когда нужна коррекция.
Глубокое погружение
Окна и watermark: работа с event time
Клик случается в момент — его event time (время события), проставленное на устройстве. Но приходит он на сервер позже — processing time (время обработки) — и зазор дико переменчив: пользователь в метро может выдать клик, чьё событие придёт через 40 минут. Если агрегировать по processing time («считай, что показалось в эту минуту»), поздние клики падают не в ту минуту, и твой поминутный биллинг неверен. Поэтому агрегируешь по event time, группируя клики в окна по тому, когда они случились, а не когда пришли.
Это поднимает проблему закрытия: окно 14:05:00–14:05:59 не может ждать опоздавших вечно, но мгновенное закрытие выбрасывает поздних. Механизм — watermark — утверждение движка «я считаю, что уже видел все события до event-time T». Когда watermark проходит конец окна, окно закрывается и выдаёт счёт. Watermark намеренно отстаёт от последнего event time (скажем, на пару минут), чтобы поглотить обычное опоздание, меняя задержку на полноту.
окно [14:05:00, 14:06:00)
клик A event-time 14:05:10 пришёл 14:05:12 ✅ в окне
клик B event-time 14:05:40 пришёл 14:05:43 ✅ в окне
watermark достигает 14:06:00+отставание → окно ЗАКРЫВАЕТСЯ, выдаёт счёт = 2
клик C event-time 14:05:55 пришёл 14:09:30 ⚠️ ПОЗДНО — окно уже закрыто▸Почему это работает
Почему просто не ждать дольше перед закрытием каждого окна, чтобы ничто никогда не опаздывало? Потому что отставание watermark — прямой налог на задержку каждого результата, платимый за поимку редкого опоздавшего. Если 99,9% кликов приходят за 10 секунд, но 0,1% берут 30 минут, выставление отставания watermark в 30 минут делает каждое число дашборда устаревшим на 30 минут ради спасения одного клика из тысячи. Senior-ход — размерить watermark под обычный случай (скажем, минута-две), принять, что по-настоящему поздние события минуют своё окно, и обработать их отдельным механизмом — side output с allowed-lateness или, лучше, пакетным слоем сверки, перечитывающим всё без давления времени. Ты не делаешь быстрый путь медленным ради корректности; ты делаешь быстрый путь быстрым и даёшь медленному пути его поправить.
Exactly-once, идемпотентность и дедуп
«Exactly-once» — заглавная гарантия, и она тончайшая. И сеть, и конвейер создают дубликаты: клиент ловит таймаут и ретраит тот же клик (событие приходит дважды с тем же click_id); потоковый процессор падает, посчитав событие, но до коммита оффсета, затем реплеит его. Наивно оба раздувают счёт — пересчёт из вступления.
Починка имеет две половины, для двух источников дублирования:
- Дедуп на источнике по ключу идемпотентности. Каждый клик несёт уникальный
click_id, сгенерированный на краю. Агрегатор дедупит по нему — внутри окна уже виденныйclick_idотбрасывается. Это убивает клиентские/сетевые ретраи. (Состояние дедупа само ограничено окном: нужно помнить ID лишь для ещё открытого окна.) - Сквозная exactly-once обработка. Это не «событие доставлено один раз» (невозможно по ненадёжной сети); это «эффект на результат применён один раз». Современные потоковые движки достигают этого, связывая коммит выхода и продвижение входного оффсета в одну атомарную воспроизводимую единицу — Flink делает это периодическими checkpoint’ами (консистентный снимок всего состояния операторов плюс потреблённые оффсеты) и транзакционным sink’ом (транзакции Kafka или идемпотентная запись по ключу окна), так что при восстановлении движок отматывается к последнему checkpoint и переизлучает без двойного применения.
Ментальная модель: exactly-once — это at-least-once доставка плюс идемпотентное/транзакционное применение. Ты не остановишь приход дубликатов; ты делаешь дублирующее применение невозможным.
▸Частая ошибка
Классическая ошибка — верить, что флаг «exactly-once» от вендора означает физическую доставку сообщений ровно один раз — а затем класть неидемпотентный сайд-эффект (списать с карты, отправить письмо, инкрементить сырой внешний счётчик) прямо в потоковый оператор. Когда задача восстанавливается из checkpoint и реплеит, этот сайд-эффект срабатывает снова, потому что exactly-once покрывает лишь состояние внутри границы checkpoint/транзакции движка. Списание с карты — вне её. Починка — держать внешние сайд-эффекты идемпотентными сами по себе (ключи идемпотентности на списании) или толкнуть их за границу в транзакционный sink, контролируемый движком. Трактуй «exactly-once» как свойство внутреннего состояния движка и его закоммиченного выхода, никогда как обещание о произвольных эффектах, что ты дёргаешь посреди потока.
Сверка: пакетный слой проверяет поток
Потоковый слой быстр, но приближён: он мог закрыть окно до прихода позднего клика или коротко посчитать дважды на грани восстановления после крэша. Поэтому задача пакетной сверки гоняется (ежечасно/ежедневно) по полному набору сырых событий в холодном хранилище — без давления watermark, каждый поздний клик теперь на месте — и считает авторитетный итог. Она сравнивает с потоковым итогом и поправляет цифру биллинга. Это страховочная сеть, которой не хватало команде из вступления: дашборд может быть чуть неверен и самоисцелиться за минуты, но число денег приходит от пакетной задачи, видевшей всё. Паттерн обобщается: пусть быстрый путь будет eventually-consistent, а медленный путь — источником истины для всего биллингуемого.
Узкие места и компромиссы
Ядро компромисса — задержка против полноты, воплощённое в watermark: длиннее отставание ловит больше поздних событий, но устаревание растёт, так что тюнишь под обычный случай и даёшь пакету разобраться с хвостом. Горячие партиции кусают на уровне лога — если партиционируешь по ad_id и одна вирусная реклама берёт 30% кликов, потребитель той партиции перегружен; добавляешь соль/подключ, чтобы размазать горячий ключ по партициям, затем переагрегируешь. Размер состояния — ограничение потокового движка: оконные множества дедупа и состояние агрегации должны влезать в управляемое состояние (на RocksDB, в checkpoint’ах) — ограничивай его, закрывая окна вовремя и истекая ключи дедупа. Цена lambda — поддержка двух кодовых баз (поток + пакет), которые должны совпадать, и ровно поэтому kappa (один движок, реплей для коррекции) привлекательна — ценой более тяжёлой потоковой системы и долгих реплеев. И под всем этим у exactly-once есть цена пропускной способности: checkpoint’ы и транзакции добавляют координацию, так что система, терпящая редкий двойной счёт (чистая аналитика), часто гоняет at-least-once ради скорости — но биллинг не может, поэтому здесь платишь налог.
Событие клика приходит через 30 минут после того, как случилось, потому что пользователь был на дрянной связи. Твои поминутные окна агрегируют по event time с отставанием watermark в 2 минуты. Что происходит с этим кликом, и как его всё-таки выставить в счёт?
Твоя потоковая задача настроена на «exactly-once», но ты поставил HTTP-вызов `chargeAdvertiser()` прямо внутрь оператора агрегации. После крэша задача восстанавливается из последнего checkpoint и реплеит — и часть рекламодателей списаны дважды. Почему?
Практическое exactly-once — это на самом деле at-least-once доставка плюс _______ применение: ты не остановишь приход дубликатов по ненадёжной сети, поэтому вместо этого делаешь так, чтобы применить одно событие дважды не давало лишнего эффекта — через ключ дедупа или транзакционный коммит, связанный с входным оффсетом.
- 01Почему агрегировать по event time, и какова роль watermark?
- 02Что на самом деле значит «exactly-once» здесь и как достигается?
- 03Для чего пакетный слой сверки и как он связан с lambda/kappa?
Подсчёт кликов по рекламе для биллинга — это задача корректности под маской задачи подсчёта, потому что сеть доставляет события не по порядку, поздно и иногда дважды, а процессоры падают на полпути — сплетённый пересчёт из вступления. Архитектура — lambda над durable, воспроизводимым логом: потоковый слой скорости даёт счётчики дашборда свежестью в секунды, а пакетный слой перечитывает неизменяемые сырые события, чтобы посчитать авторитетный итог биллинга и свести мелкие ошибки быстрого пути (kappa — одно-движковая альтернатива с реплеем-для-коррекции). Быстрый путь агрегирует по event time, чтобы поздние клики не падали в неверную минуту, закрывая окна по watermark, чьё отставание размерено под обычный случай — принимая, что редкие опоздавшие минуют окно и восстанавливаются пакетом. Exactly-once инженерится как at-least-once доставка плюс идемпотентное/транзакционное применение: дедуп ретраев по click_id, связка коммита выхода с входным оффсетом через checkpoint’ы и транзакционный sink — но он покрывает лишь состояние внутри движка, так что внешние сайд-эффекты надо делать идемпотентными самими. Постоянные компромиссы — задержка против полноты (watermark), горячие партиции на вирусной рекламе (соль на ключ), ограниченное состояние потока (закрывай окна, истекай ключи дедупа) и налог на пропускную способность exactly-once, который платишь, потому что, в отличие от аналитики, биллинг не терпит ни единого двойного счёта. Теперь, когда увидишь потоковый оператор с вызовом внешнего сервиса внутри, сразу спросишь: что произойдёт при реплее из чекпойнта? Если ответ «этот вызов сработает снова» — перенеси эффект за транзакционную границу или сделай его идемпотентным в своём праве.
Практика
Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.