Skip to content
12 / 15

Стриминговая обработка: окна, состояние и опоздавшие события

Стриминговая обработка отличается от пакетной тем, что данные бесконечны, поэтому агрегаты считают по окнам: фиксированным (tumbling), скользящим с перекрытием (hopping), сессионным (разрыв активности) или скользящим по каждому событию (sliding). Ключевая сложность — время: время события (когда произошло) и время обработки (когда дошло) не совпадают, а события приходят с задержкой и не по порядку. Отсюда понятие watermark — оценка «событий старше этого момента больше не ждём»; она задаёт, когда окно можно закрыть и выдать результат. Опоздавшие после watermark события либо отбрасываются, либо обновляют уже выданный результат, либо уходят в отдельный поток — выбор диктует бизнес. Второй тяжёлый вопрос — состояние: агрегаты и соединения требуют хранить данные между событиями, и это состояние нужно восстанавливать после сбоя, обычно через журнал изменений в брокере и контрольные точки.

Стриминговая обработка: окна, состояние и опоздавшие события | JScriptiser