Сообщения

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

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

Каждое звено обещает своё, а отвечать приходится за цепочку целиком: как 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

Flink SQL, коннекторы и эксплуатация: от Table API до прода

Flink SQL, коннекторы и эксплуатация: от Table API до прода

Писать Flink-джобы на низкоуровневом API нужно не всегда: Table API и Flink SQL закрывают большую часть задач декларативно, включая стриминговые джойны и агрегации. Плюс то, без чего не выйти в прод: коннекторы (Kafka, CDC, lakehouse), управление параллелизмом, диагностика backpressure, тюнинг чекпоинтов и апгрейд джобы через savepoint без потери состояния. Завершает погружение практикой эксплуатации

Состояние и 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

MirrorMaker 2 и геораспределённый Kafka: репликация между кластерами, DR, active-active

MirrorMaker 2 и геораспределённый Kafka: репликация между кластерами, DR, active-active

Один кластер Kafka живёт в одном дата-центре — а бизнесу нужны катастрофоустойчивость и присутствие в нескольких регионах. Разбираем два подхода: растянутый (stretch) кластер против репликации отдельных кластеров через MirrorMaker 2; что MM2 на самом деле переносит (топики, offset’ы, ACL, конфиги) и как работает трансляция offset при failover; топологии active-passive/active-active/hub-and-spoke; RPO/RTO, failover и failback; цена латентности и трафика между регионами. С живым стендом из двух KRaft-кластеров и MM2.

Kafka-клиенты в разных языках: драйверы, родословные, различия

Kafka-клиенты в разных языках: драйверы, родословные, различия

Клиент Kafka делает куда больше, чем «подключиться»: партиционирование, батчинг, сжатие, управление offset, ребаланс, идемпотентность, транзакции — и всё это разные библиотеки реализуют по-разному. Три родословные (референсный JVM-клиент, обёртки над librdkafka на C, чистые реализации протокола), почему у Go сразу несколько драйверов (franz-go/sarama/segmentio/confluent) и чем они отличаются, и что из этого следует для фич, производительности и деплоя. С живыми примерами produce/consume на Java, Go, Rust, Python, C#, C++, Scala

В работе

Ближайшие материалы, которые продолжают этот раздел.

Сообщения

RabbitMQ, Kafka, NATS: как не выбирать брокер по моде

Вводный материал про три разных класса messaging-решений и то, какие задачи они на самом деле закрывают.

Сообщения

NATS для микросервисов и внутренних интеграций

Когда NATS оказывается достаточно и почему это не просто "Kafka полегче".

Сообщения

Kafka в Docker Compose для разработки

Как поднять понятный dev-стенд Kafka без иллюзий, что это уже production.