Основы потоковой обработки
Стриминг обрабатывает события по мере поступления, с латентностью в секунды и миллисекунды. Собес проверяет, выбираешь ли ты его по требованию к свежести, а не по моде, и понимаешь ли, что лог это и транспорт, и буфер.
Стержень: батч проще и дешевле в сопровождении; стриминг берут, когда продукту реально нужна свежесть, а не «чтобы модно».
// Формулировки: «когда стриминг, а когда батч?», «что такое consumer lag?», «что такое backpressure?».
Лог и декаплинг
Основа стриминга - лог: append-only последовательность событий (Kafka-класс) как транспорт и буфер. Продюсеры пишут, консюмеры читают в своём темпе со своих offset'ов, независимо друг от друга.
Декаплинг - главный дар брокера: источник и потребитель не знают друг о друге. Потребитель упал - лог подождёт в пределах retention, а новый потребитель может перечитать историю с начала.
// Отсюда жёсткое требование к retention: если он короче времени простоя потребителя, события уйдут безвозвратно. Срок хранения выбирают с запасом на даунтайм и пересчёты.
- append-only log
- неизменяемая последовательность событий со смещениями
- offset
- позиция потребителя в логе партиции
События, backpressure, lag
Событие - иммутабельный факт «что произошло» со временем; состояние строится из событий свёрткой потока, а не наоборот. Исправление - не правка события задним числом (лог иммутабелен), а новое событие-компенсация.
Backpressure: потребитель медленнее продюсера - очередь растёт. Стратегии - масштабировать консюмеров, буферизовать, деградировать; игнор ведёт к OOM (out of memory) или вечному отставанию.
// Consumer lag - метрика номер один: отставание чтения от конца лога. Растущий лаг значит, что потребитель не справляется, и это видно раньше любых жалоб - если за ним следить.
- consumer lag
- отставание потребителя от конца лога
- backpressure
- давление на потребителя, не успевающего за потоком
Stateless, stateful, и не вместо батча
Обработка бывает stateless и stateful. Фильтр и маппинг тривиальны - каждое событие само по себе. А агрегаты, джойны и окна требуют состояния, а состояние требует чекпоинтов и восстановления после сбоя.
Стриминг не заменяет батч: типовая архитектура держит поток для operational-свежести плюс батч-пересчёт для точности и исправлений (или lakehouse с инкрементами).
// Поэтому стриминг ради дашборда, который смотрят раз в день, переинженерия в 10 раз дороже без пользы: сложность реального времени при отсутствии потребителя этой скорости.
- stateless / stateful
- без состояния (фильтр) / с состоянием (окна, агрегаты)
- event-driven
- состояние как свёртка потока иммутабельных событий
Как отвечать: «Когда стриминг оправдан, а когда батч?»
Выбираю по требованию потребителя к свежести, а не по хайпу. Если решение принимается на данных возрастом в секунды - антифрод, мониторинг, оперативные алерты, нужен стриминг. Если потребителю хватает «раз в час» или «раз в день» - беру батч, он проще, дешевле и надёжнее в сопровождении. Стриминг платит постоянной сложностью: состояние с чекпоинтами, обработка поздних событий, мониторинг lag и backpressure, и это оправдано, только когда свежесть реально ценна. Классическая ошибка - поднять стриминг ради отчёта, который открывают раз в день. Причём одно другого не исключает: типовая зрелая архитектура - поток для оперативной свежести плюс батч-пересчёт для точности и исправлений, потому что поток даёт «лучшее знание на сейчас», а батч задним числом всё аккуратно сводит.
Почему это сильный ответ: критерий - свежесть для потребителя, названа постоянная цена стриминга, и упомянута гибридная архитектура (поток + батч-пересчёт) - зрелый взгляд, а не «стриминг круче».
На чём валят
- −Стриминг для дашборда, который смотрят раз в день - цена x10 без пользы.
- −Не мониторить lag - «данные отстали на 6 часов» узнаётся от бизнеса.
- −Retention короче, чем даунтайм потребителя - события ушли безвозвратно.
- −Мутировать события задним числом - лог иммутабелен, исправление = событие-компенсация.
- −Состояние в памяти без чекпоинтов: рестарт = агрегаты с нуля.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 14, остальные разбираются в тренажёре.
- Что даёт consumer group в Kafka?A)Строгую гарантию, что каждое сообщение прочитают сразу все консьюмеры группы одновременноB)Распределение партиций топика между консьюмерами группы для параллельного чтенияC)Автоматическое сжатие всех сообщений топика для экономии места на дискеD)Полную остановку топика на время обработки одного сообщения
показать ответ и разбор
+B)Распределение партиций топика между консьюмерами группы для параллельного чтения// разбор: В consumer group партиции топика распределяются между консьюмерами: каждую партицию читает ровно один консьюмер группы, что даёт параллелизм чтения (но не больше, чем число партиций). При падении или добавлении консьюмера происходит ребалансировка — партиции переназначаются. Разные группы читают один топик независимо, каждая со своими offset'ами.
- Почему Kafka гарантирует порядок сообщений только внутри партиции, а не по всему топику?A)Kafka вообще не обеспечивает никакого порядка сообщений в этом случаеB)Порядок теряется из-за сжатия сообщений при их записи на диск брокераC)Система Kafka обеспечивает строгий полный глобальный порядок по всему топику целикомD)Партиции читаются независимо и параллельно — общий порядок между ними не определён
показать ответ и разбор
+D)Партиции читаются независимо и параллельно — общий порядок между ними не определён// разбор: Каждая партиция — отдельный упорядоченный лог, но партиции обрабатываются независимыми консьюмерами параллельно, поэтому единого порядка между сообщениями разных партиций нет. Чтобы связанные события (например, все действия одного пользователя) шли строго по порядку, их направляют в одну партицию через ключ сообщения (partition key).
- Consumer lag растёт: консьюмер не успевает за продюсером. Что это значит и что делать?A)Обработка медленнее притока; масштабировать консьюмеров/партиции, оптимизироватьB)Растущий consumer lag означает, что брокер Kafka уже потерял все сообщения топикаC)Растущий lag — это нормально и не требует реакцииD)Достаточно просто перезапустить продюсера, и lag немедленно исчезнет сам собой
показать ответ и разбор
+A)Обработка медленнее притока; масштабировать консьюмеров/партиции, оптимизировать// разбор: Consumer lag — разрыв между последним записанным offset'ом и обработанным консьюмером. Устойчивый рост значит, что обработка медленнее притока: растёт задержка, а при исчерпании ретенции возможна и потеря необработанных сообщений. Лечат масштабированием консьюмеров (но не больше числа партиций), оптимизацией обработки, батчингом и буферизацией. Параллелизм консьюмеров упирается в число партиций.
- Что такое продюсер и консьюмер в системе вроде Kafka?A)Продюсер пишет сообщения в топик, консьюмер их читает; они развязаны через брокерB)Продюсер и консьюмер должны быть постоянно соединены напрямую и работать строго синхронноC)Продюсер читает данные из топика, а консьюмер записывает их тудаD)Это два названия одного процесса, который одновременно и пишет, и читает одни и те же данные
показать ответ и разбор
+A)Продюсер пишет сообщения в топик, консьюмер их читает; они развязаны через брокер// разбор: Продюсер публикует сообщения в топик, консьюмер подписывается и читает их. Они не знают друг о друге и работают в своём темпе — брокер (Kafka) хранит сообщения между ними как буфер/лог. Эта развязка (decoupling) — суть шины сообщений: продюсер не ждёт консьюмера, консьюмеров может быть много и они читают независимо, а пики сглаживаются очередью.
- Почему у стриминга обычно ниже задержка, чем у батча, но сложнее корректность?A)Событие за событием → низкий лаг, но поздние данные и состояние усложняютB)Стрим одновременно и быстрее, и проще батча во всех без исключения аспектах обработкиC)Стрим медленнее батча, зато точен и не знает проблемы поздних данныхD)Разница в объёме данных: стрим применяют к маленьким наборам данных
показать ответ и разбор
+A)Событие за событием → низкий лаг, но поздние данные и состояние усложняют// разбор: Батч копит данные за интервал и обрабатывает разом: просто и легко пересчитать, но результат запаздывает на длину окна. Стрим реагирует на каждое событие — задержка секунды/миллисекунды, но за это платит сложностью: приходится управлять состоянием на лету, разбираться с поздними и внеочередными событиями, гарантиями доставки и восстановлением после сбоя без полного пересчёта.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.