Потоковая обработка данных
Как поверх at-least-once избежать повторного эффекта?
At-least-once допускает повтор после сбоя между эффектом и подтверждением. Повторный эффект убирают устойчивым event ID и атомарным INSERT ... ON CONFLICT, dedup-таблицей либо идемпотентным upsert. Offset или checkpoint фиксируют вместе с результатом, если sink поддерживает общую транзакцию; иначе нужна согласованная recovery-схема. Kafka transactions дают exactly-once для Kafka-to-Kafka, но не делают произвольный внешний API транзакционным. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.
Ссылки для изучения
Как Kafka распределяет партиции внутри consumer group?
Внутри одной consumer group каждая partition в данный момент назначена не более чем одному consumer, а один consumer может читать несколько partitions. Поэтому параллелизм группы ограничен числом partitions; лишние consumers простаивают. Порядок Kafka гарантирует только внутри partition. Вступление, выход или сбой участника вызывает rebalance, поэтому обработчик должен корректно завершить работу, зафиксировать offset и пережить повторную доставку. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.
Ссылки для изучения
Примеры хороших ответов из реальных собеседований
- Java мок-интервью: сервис заказов, Kafka и PostgreSQL · 21:47–22:41Мок-собеседование · Совместный разбор
На двух примерах разбирают распределение партиций в одной consumer group: лишний consumer простаивает, а при меньшем числе consumers одному из них достаётся несколько партиций.
Как устроены event-time окна и watermark?
Tumbling windows не пересекаются, sliding могут перекрываться, session объединяют события до разрыва активности. Окна считают по event time из события, а watermark выражает оценку, насколько далеко поток продвинулся, и позволяет закрывать окна и очищать state. Allowed lateness определяет обработку опоздавших событий до этой границы. События за watermark могут быть отброшены, обновлены или отправлены отдельно, в зависимости от engine и output mode.
Ссылки для изучения
Что ограничивает join двух потоков?
Stream-stream join требует ключа и конечной временной связи, например событие B в пределах часа от A. Без границы движок вынужден бесконечно хранить прошлые строки. Event-time constraints и watermarks позволяют удалить state, когда совпадение уже считается невозможным. Цена зависит от cardinality ключей, skew и lateness. Очень поздние события могут не соединиться, поэтому нужен отдельный путь или более широкий state с большей стоимостью.
Ссылки для изучения
Чем micro-batch отличается от record-at-a-time streaming?
Micro-batch накапливает события за короткий интервал и запускает пакетный план, обычно получая высокий throughput и задержку не ниже длительности триггера и обработки. Record-at-a-time engine продвигает каждое событие непрерывно и может дать меньшую latency, но имеет другую стоимость планирования и checkpoint. Гарантии при сбое определяются источником, state и sink, а не одним названием модели. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.
Ссылки для изучения
Что нужно Kafka для безопасной работы в Kubernetes?
Broker требует стабильную identity, persistent volume и корректные advertised listeners; обычно это оформляет operator или StatefulSet-подобное управление. Реплики разносят anti-affinity по nodes и zones, задают disruption budget и аккуратный rolling update. Requests и limits учитывают page cache и disk throughput. Мониторят ISR, under-replicated partitions, controller, disk, request latency и consumer lag; один Kubernetes restart не должен уничтожать единственную копию данных.
Ссылки для изучения





