Streaming

Streams, Lua и функции: программируемость Redis

Streams, Lua и функции: программируемость Redis

Redis умеет не только хранить, но и выполнять логику на сервере. Разбираем три способа программируемости: Streams с consumer groups (журнал сообщений с подтверждениями), Lua-скрипты и Redis Functions (атомарные серверные процедуры), плюс обзор модулей. Честно проводим границу Streams vs Kafka — где Redis Streams достаточно, а где нужна полноценная брокерная платформа.

Сквозные гарантии и жизнь пайплайна

Сквозные гарантии и жизнь пайплайна

Каждое звено обещает своё, а отвечать приходится за цепочку целиком: как at-least-once складывается в effectively-once, где заканчивается транзакция Kafka, почему дедуп на стоке обязателен и как пайплайн переживает replay, эволюцию схемы и битые события

Собираем пайплайн: PostgreSQL → Debezium → Kafka → ClickHouse

Собираем пайплайн: PostgreSQL → Debezium → Kafka → ClickHouse

Сквозная цепочка от изменения в PostgreSQL до строки в витрине: Debezium как коннектор Kafka Connect, outbox event router, snapshot и переход в streaming, сток в ClickHouse. Не про механику WAL и не про внутренности Kafka — про то, как это собирается вместе и где проходят швы

Flink DataStream API: когда SQL мало

Flink DataStream API: когда SQL мало

Flink SQL/Table API и DataStream — два API над общим runtime, и для оконных агрегатов с джойнами декларативного хватает. Но реакция по таймеру, произвольное состояние, корреляция двух потоков с кастомной логикой, обогащение через broadcast/async и маршрутизация в side output живут в Java-первом DataStream API. Где кончается SQL и что даёт спуск: KeyedProcessFunction, состояние руками, event-time vs processing-time таймеры, connect + CoProcessFunction

Состояние и exactly-once в Flink: checkpoints, savepoints, 2PC

Состояние и exactly-once в Flink: checkpoints, savepoints, 2PC

Стрим-процессор без надёжного состояния бесполезен: агрегаты, джойны, дедуп — всё это состояние, которое нельзя потерять при падении. Как Flink это решает: keyed vs operator state, state backends (heap vs RocksDB для состояния больше памяти), распределённые снапшоты через барьеры (алгоритм Chandy-Lamport), инкрементальные checkpoints и savepoints для апгрейда, и exactly-once end-to-end через two-phase commit в стоки

Модель Flink: dataflow, event-time, watermarks и окна

Модель Flink: dataflow, event-time, watermarks и окна

Что делает Flink настоящим стрим-процессором, а не «циклом по сообщениям»: dataflow-граф операторов на JobManager/TaskManager, различие event-time и processing-time (и почему первое — единственно честное для аналитики потоков), watermarks как механизм «мы уже видели всё до момента T», окна (tumbling/sliding/session) и что делать с опоздавшими данными. Фундамент серии про Flink

Kafka: лог, топики, партиции — ментальная модель

Kafka: лог, топики, партиции — ментальная модель

Kafka — это распределённый лог, а не очередь: топики и партиции, offset, producer и consumer, как ключ определяет партицию и порядок, зачем репликация. Вводная серии «Kafka вглубь»: правильная ментальная модель и когда Kafka уместен, а когда нет

Streams в RabbitMQ: Kafka-подобный лог и можно ли заменить Kafka

Streams в RabbitMQ: Kafka-подобный лог и можно ли заменить Kafka

Что такое streams в RabbitMQ — append-only лог с offset-чтением и репликацией, как их объявлять и читать (native stream-протокол и AMQP), super streams для партиционирования, честное сравнение с Kafka и разбор, где RabbitMQ Streams реально заменяет Kafka, а где нет