---
title: Потоковая обработка данных
questionDates:
  q-transfer-0305: '2026-10-01'
  q-transfer-0306: '2026-10-01'
  q-transfer-0307: '2026-10-01'
  q-transfer-0308: '2026-10-01'
  q-transfer-0309: '2026-10-01'
  q-transfer-0325: '2026-10-01'
seo:
  description: >-
    Тема «Потоковая обработка данных» для собеседования Data Engineer. Как
    поверх at-least-once избежать повторного эффекта? Как Kafka распределяет
    партиции внутри consumer group?
  title: Потоковая обработка данных — Data Engineer
---

[Все темы Data Engineer](/prep/data-engineer)

## Как поверх at-least-once избежать повторного эффекта? [#q-transfer-0305]

At-least-once допускает повтор после сбоя между эффектом и подтверждением. Повторный эффект убирают устойчивым event ID и атомарным `INSERT ... ON CONFLICT`, dedup-таблицей либо идемпотентным upsert. Offset или checkpoint фиксируют вместе с результатом, если sink поддерживает общую транзакцию; иначе нужна согласованная recovery-схема. Kafka transactions дают exactly-once для Kafka-to-Kafka, но не делают произвольный внешний API транзакционным. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.

:::note[Ссылки для изучения]

1. [Повторная доставка и границы exactly-once: Kafka-to-Kafka и внешние системы](https://habr.com/ru/companies/ydb/articles/972180/)
   :::

---

## Как Kafka распределяет партиции внутри consumer group? [#q-transfer-0306]

Внутри одной consumer group каждая partition в данный момент назначена не более чем одному consumer, а один consumer может читать несколько partitions. Поэтому параллелизм группы ограничен числом partitions; лишние consumers простаивают. Порядок Kafka гарантирует только внутри partition. Вступление, выход или сбой участника вызывает rebalance, поэтому обработчик должен корректно завершить работу, зафиксировать offset и пережить повторную доставку. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.

:::note[Ссылки для изучения]

1. [Kafka consumer group: назначение разделов, offsets и ребалансировка](https://yandex.cloud/ru/docs/managed-kafka/concepts/producers-consumers#consumer-groups)
   :::

---

## Как устроены event-time окна и watermark? [#q-transfer-0307]

Tumbling windows не пересекаются, sliding могут перекрываться, session объединяют события до разрыва активности. Окна считают по event time из события, а watermark выражает оценку, насколько далеко поток продвинулся, и позволяет закрывать окна и очищать state. Allowed lateness определяет обработку опоздавших событий до этой границы. События за watermark могут быть отброшены, обновлены или отправлены отдельно, в зависимости от engine и output mode.

:::note[Ссылки для изучения]

1. [Structured Streaming: tumbling/sliding/session windows и watermark](https://learn.microsoft.com/ru-ru/fabric/data-engineering/structured-streaming-stateful-processing)
   :::

---

## Что ограничивает join двух потоков? [#q-transfer-0308]

Stream-stream join требует ключа и конечной временной связи, например событие B в пределах часа от A. Без границы движок вынужден бесконечно хранить прошлые строки. Event-time constraints и watermarks позволяют удалить state, когда совпадение уже считается невозможным. Цена зависит от cardinality ключей, skew и lateness. Очень поздние события могут не соединиться, поэтому нужен отдельный путь или более широкий state с большей стоимостью.

:::note[Ссылки для изучения]

1. [Watermark двух потоков: ограничение состояния, поздние данные и min/max policy](https://learn.microsoft.com/ru-ru/azure/databricks/structured-streaming/watermarks)
2. [JOIN: начало разбора. karpov.courses, с 0:08](https://www.youtube.com/watch?v=Xy3RWYKRVb4&t=8s)

:::

---

## Чем micro-batch отличается от record-at-a-time streaming? [#q-transfer-0309]

Micro-batch накапливает события за короткий интервал и запускает пакетный план, обычно получая высокий throughput и задержку не ниже длительности триггера и обработки. Record-at-a-time engine продвигает каждое событие непрерывно и может дать меньшую latency, но имеет другую стоимость планирования и checkpoint. Гарантии при сбое определяются источником, state и sink, а не одним названием модели. Гарантию формулируют вместе с границей состояния, источником, sink и поведением при восстановлении после сбоя.

:::note[Ссылки для изучения]

1. [Микропакетная и непрерывная обработка в Azure Databricks: задержка, ресурсы и checkpoints](https://learn.microsoft.com/ru-ru/azure/databricks/structured-streaming/real-time)
   :::

---

## Что нужно Kafka для безопасной работы в Kubernetes? [#q-transfer-0325]

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 не должен уничтожать единственную копию данных.

:::note[Ссылки для изучения]

1. [Состояние и постоянное хранилище в Kubernetes: StatefulSet, PV и оператор Kafka Strimzi](https://habr.com/ru/companies/flant/articles/809377/)
   :::
