A B C D E F G H I J K L M N O P Q R S T V W Y Z А Б В Г Е И К М О П С Т Ц

Event streaming

Event streaming

 

Event streaming (потоковая передача событий) это модель обмена данными, при которой каждое значимое изменение в системе записывается в неизменяемый упорядоченный лог, хранится там заданное время и может быть прочитано любым числом независимых потребителей. Событием считается факт, который уже случился, например заказ оформлен, платёж прошёл, датчик выдал измерение. Ключевое отличие от привычной интеграции через очередь в том, что чтение не уничтожает запись. Лог остаётся на диске, потребители лишь двигают по нему свою позицию. Из одного этого свойства вырастают и повторное чтение истории, и независимые группы потребителей, и вся событийно-ориентированная архитектура.

 

Что такое Event Streaming и какую задачу он закрывает

Классическая интеграция сервисов строится на прямых вызовах. Сервис заказов дёргает сервис доставки, оба ждут друг друга. Пока сервисов пять, схема работает. Когда их пятьдесят, а на каждое событие подписаны аналитика, биллинг, антифрод и рекомендации, число связей растёт быстрее, чем команда успевает их поддерживать.

Потоковая передача событий разворачивает картину. Производитель данных не знает, кто их прочитает, и пишет событие в топик один раз. Потребители подписываются сами и обрабатывают поток в своём темпе. Отправитель не ждёт получателя, получатель не блокирует отправителя, а новый подписчик добавляется без единой строчки изменений на стороне источника.

Стоит развести три понятия, которые в разговоре часто сваливают в кучу. Обработка в реальном времени (real-time processing) это про задержку. Потоковая обработка (stream processing) это про вычисления над непрерывным потоком, чем занимаются Apache Flink или Kafka Streams. Event streaming это про сам транспорт и хранение, а именно про то, как события попадают в лог и как оттуда достаются. Одно другому не мешает, но и не заменяет.

 

Архитектура лога, партиции и оффсеты

Платформа потоковой передачи событий устроена вокруг одной структуры данных, а именно журнала с добавлением в конец. Разобравшись, как он работает, вы поймёте и правила порядка событий, и то, почему повторное чтение истории ничего не ломает.

Топик из трёх партиций, оффсеты записей и две независимые группы потребителей, читающие один и тот же лог

 

 

 

Топик как долговременный журнал

Топик (topic) это именованный поток событий, физически разложенный на партиции (partitions). Партиция представляет собой файл, куда записи добавляются строго в конец и больше никогда не меняются. Позиция записи внутри партиции называется оффсетом (offset) и растёт монотонно с нуля. Порядок гарантирован в пределах партиции, а не топика целиком, и это первое, обо что спотыкаются новички.

Партицию для события выбирает хеш ключа. Все события одного ключа неизбежно попадают в одну партицию, поэтому история конкретного заказа или конкретного магазина читается в правильной последовательности. События без ключа раскладываются по партициям равномерно, и тогда порядок между ними не определён вообще.

Срок жизни записи задаётся политикой хранения, а не фактом прочтения. На нашем стенде параметр KAFKA_LOG_RETENTION_HOURS выставлен в 168, то есть события лежат в логе неделю независимо от того, забрал их кто-нибудь или нет.

 

Продюсеры, потребители и группы

Продюсер (producer) пишет события в топик и получает подтверждение с номером партиции и оффсетом. Потребитель (consumer) читает лог и сам хранит на брокере свою позицию. Между ними стоит группа потребителей (consumer group), без которой масштабирование чтения не собирается.

Правило простое. Партиция достаётся ровно одному потребителю внутри группы, а разные группы читают один и тот же лог полностью и независимо друг от друга. Отсюда две вещи сразу.

  • Масштабирование внутри группы. Добавили второй экземпляр сервиса, партиции переразделились между экземплярами, пропускная способность выросла. Потолок это число партиций, лишние потребители простаивают.
  • Независимость между группами. Аналитика и биллинг читают одни и те же события, не зная друг о друге и не мешая друг другу. Тормозящая группа не задерживает остальных.

Именно связка партиций и групп даёт горизонтальное масштабирование без изменения кода производителя.

 

Принцип работы, гарантии доставки и повторное чтение

Продюсер получает подтверждение записи в зависимости от настройки acks. При значении all брокер отвечает только после того, как запись попала во все синхронные реплики, поэтому потеря одного узла не уносит событие с собой. На нашем стенде запись восьми событий с таким режимом заняла 0.024 секунды, то есть сам лог тут вообще не узкое место.

Гарантий доставки три. Режим at most once допускает потерю, at least once допускает дубли, exactly once убирает и то и другое ценой транзакций и заметного усложнения. В реальных конвейерах чаще всего берут at least once и делают обработчик идемпотентным, потому что дубль дешевле потери.

Отдельно стоит повторное чтение, оно же replay. Поскольку события никуда не деваются после обработки, новая группа потребителей может подключиться и поднять всю историю с нулевого оффсета. Так чинят баг в логике обработки, наполняют данными новый сервис, пересчитывают витрину по изменившейся формуле. Боевые потребители при этом ничего не замечают, ведь у новой группы свои оффсеты.

Apache Kafka для инженеров данных

Код курса
DEVKI
Ближайшая дата курса
12 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Event Streaming, очередь сообщений и Event Sourcing

Три модели постоянно путают, хотя задачи у них разные. Очередь сообщений (message queue) вроде RabbitMQ доставляет задание исполнителю и удаляет его после подтверждения. Event sourcing это способ хранить состояние приложения как последовательность изменений вместо текущего снимка. Разница между двумя последними подходами разобрана в статье Event Streaming vs Event Sourcing, а сводка по трём выглядит так.

Сравнение очереди сообщений, где запись исчезает после подтверждения, и лога событий с сохранением истории

 

 

Свойство Очередь сообщений Потоковая передача событий Event Sourcing
Судьба записи удаляется после подтверждения живёт по политике хранения хранится вечно
Число получателей обычно один исполнитель любое число независимых групп агрегат восстанавливается из своего потока
Повторное чтение невозможно штатный сценарий основной механизм работы
Что передаётся команда, задание к выполнению факт, который уже случился изменение состояния объекта
Область применения фоновые задачи, распределение нагрузки интеграция систем, аналитика, конвейеры данных внутреннее устройство одного сервиса

Практический вывод такой. Очередь берут, когда задание нужно выполнить один раз и забыть. Стриминг берут, когда факт нужен многим и надолго. Event sourcing это решение об устройстве конкретного сервиса, и оно часто уживается со стримингом в одной системе.

 

Куда движется потоковая передача событий

Развитие идёт по линии упрощения эксплуатации. Начиная с Apache Kafka 4.0 режим KRaft стал обязательным, а внешний ZooKeeper из схемы исчез совсем. Кластер теперь состоит из брокеров и контроллеров кворума, второй системы координации поднимать не надо.

Второе направление это сближение с моделью очереди. Механизм share groups из официальной документации Apache Kafka позволяет нескольким потребителям разбирать записи одной партиции с подтверждением каждой по отдельности, то есть даёт классическую семантику очереди поверх лога. Предложение оформлено как KIP-932 и доехало до общего доступа в версии 4.2.

Третье направление самое затратное для инфраструктуры. Предложение KIP-1150 о бездисковых топиках принято 2 марта 2026 года и уносит хранение в объектное хранилище, убирая межзональный трафик репликации, который в облаке составляет заметную долю счёта.

 

Сценарии применения и ограничения

Потоковая передача событий закрывает не любую интеграцию, и трезвый взгляд на границы применимости экономит месяцы. Хорошо она ложится на несколько типовых задач.

  • Интеграция сервисов через шину событий. Источник публикует факты, подписчики появляются и исчезают без правок на его стороне.
  • Конвейеры данных. Захват изменений из баз через CDC, доставка в хранилище и озеро, наполнение витрин почти в реальном времени.
  • Аналитика на лету. Скользящие окна, агрегаты и детекция аномалий поверх потока, обычно средствами Apache Flink версии 2.3.0 или Kafka Streams.
  • Аудит и восстановление. Неизменяемый лог сам по себе является журналом того, что происходило в системе.

Общее у этих сценариев одно, а именно ценность самого факта события для нескольких потребителей сразу.

Ограничения тоже стоит держать в голове. Эксплуатация кластера требует людей, понимающих партиции, репликацию и ребалансировку. Порядок гарантируется только внутри партиции, поэтому глобальную последовательность придётся либо собирать ключами, либо принимать её отсутствие. Число партиций легко увеличить и невозможно уменьшить, а увеличение ломает привязку ключей к партициям. Наконец, лог хранит сырые факты без схемы, и без реестра схем изменение формата у производителя ломает всех потребителей разом.

Потоковая обработка данных с помощью Apache Flink

Код курса
FLINK
Ближайшая дата курса
29 сентября, 2026
Продолжительность
16 ак.часов
Стоимость обучения
51 200

 

Практика на демо-стенде Apache Kafka 4.3.0

По традиции весь код используемый в статье выкладываем на наш GitHub репозиторий
Потоковая передача событий (Event Streaming) from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime with DAG( dag_id="spark_submit_demo", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False ) as dag: run = BashOperator( task_id="run_job", bash_command="spark-submit app.py" ) GitHub code example Потоковая передача событий (Event Streaming)

Стенд поднимается одним контейнером с образом apache/kafka версии 4.3.0 в режиме KRaft.

# apache/kafka 4.3.0, KRaft без ZooKeeper, прогнано на стенде 2026-08-24
# Один брокер, он же контроллер. Для демо потоковой передачи событий этого достаточно.
services:
  kafka:
    image: apache/kafka:4.3.0
    container_name: event_streaming_kafka
    ports:
      - "9092:9092"
    environment:
      # KRaft: один узел совмещает роли брокера и контроллера
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      # реплика одна, потому что брокер один
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      # группа успевает собраться быстрее дефолтных трёх секунд
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      # события живут в логе неделю независимо от того, прочитал их кто-то или нет
      KAFKA_LOG_RETENTION_HOURS: 168
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"

Продюсер пишет восемь заказов четырёх магазинов, ключом служит идентификатор магазина. Колбэк доставки возвращает партицию и оффсет каждой записи, а в конце скрипт показывает, что события одного ключа не разъехались по разным партициям.

# confluent-kafka 2.15.0, Apache Kafka 4.3.0 в KRaft, прогнано на стенде 2026-08-24
"""Запись событий в топик: лог из трёх партиций, ключ определяет партицию."""
import json
import time
from confluent_kafka import Producer
from confluent_kafka.admin import AdminClient, NewTopic

BROKER = "localhost:9092"
TOPIC = "orders"
PARTITIONS = 3

# события четырёх магазинов, ключ это идентификатор магазина
EVENTS = [
    {"store": "msk-01", "order": 1001, "amount": 2500},
    {"store": "omsk-08", "order": 1002, "amount": 700},
    {"store": "msk-01", "order": 1003, "amount": 1900},
    {"store": "ekb-03", "order": 1004, "amount": 4300},
    {"store": "omsk-08", "order": 1005, "amount": 150},
    {"store": "kzn-05", "order": 1006, "amount": 8800},
    {"store": "msk-01", "order": 1007, "amount": 320},
    {"store": "ekb-03", "order": 1008, "amount": 1150},
]


def create_topic() -> None:
    """Топик создаётся явно: автосоздание на стенде выключено."""
    admin = AdminClient({"bootstrap.servers": BROKER})
    existing = admin.list_topics(timeout=10).topics
    if TOPIC in existing:
        print(f"топик {TOPIC} уже есть, партиций: {len(existing[TOPIC].partitions)}")
        return
    new_topic = NewTopic(TOPIC, num_partitions=PARTITIONS, replication_factor=1)
    for name, fut in admin.create_topics([new_topic]).items():
        fut.result()  # бросит исключение, если создать не удалось
        print(f"топик {name} создан, партиций: {PARTITIONS}")


def main() -> None:
    create_topic()
    # acks=all: продюсер ждёт подтверждения от всех синхронных реплик
    producer = Producer({"bootstrap.servers": BROKER, "acks": "all"})
    placed = []

    def on_delivery(err, msg):
        """Колбэк подтверждения: сюда приходит партиция и оффсет записи в логе."""
        if err is not None:
            print(f"ОШИБКА доставки: {err}")
            return
        placed.append((msg.key().decode(), msg.partition(), msg.offset()))

    t0 = time.time()
    for event in EVENTS:
        producer.produce(
            TOPIC,
            key=event["store"].encode(),
            value=json.dumps(event).encode(),
            on_delivery=on_delivery,
        )
    producer.flush(15)  # дожидаемся всех подтверждений
    elapsed = time.time() - t0

    print(f"записано событий: {len(placed)} за {elapsed:.3f} с")
    print("ключ -> партиция:оффсет")
    for key, partition, offset in placed:
        print(f"  {key} -> {partition}:{offset}")

    # события одного магазина обязаны лежать в одной партиции: порядок гарантируется в её пределах
    by_key = {}
    for key, partition, _ in placed:
        by_key.setdefault(key, set()).add(partition)
    for key, parts in sorted(by_key.items()):
        print(f"  ключ {key}: партиции {sorted(parts)}")


if __name__ == "__main__":
    main()

Прогон на стенде дал такую раскладку.

топик orders создан, партиций: 3
записано событий: 8 за 0.024 с
  msk-01 -> 0:0
  msk-01 -> 0:1
  msk-01 -> 0:2
  omsk-08 -> 1:0
  omsk-08 -> 1:1
  ekb-03 -> 2:0
  kzn-05 -> 2:1
  ekb-03 -> 2:2
  ключ ekb-03: партиции [2]
  ключ msk-01: партиции [0]
  ключ omsk-08: партиции [1]

Три события магазина msk01 легли в партицию 0 с оффсетами 0, 1 и 2 подряд, ключи omsk-08 и ekb03 ушли каждый в свою партицию. Равномерности на восьми событиях ждать не стоит, ключи для демо подобраны замером на реальном брокере. Партиционер библиотеки librdkafka работает в режиме consistent_random, и его раскладка не совпадает с murmur2 из Java-клиента, поэтому считать её на бумаге бесполезно.

Второй скрипт читает тот же лог пятью разными способами. Группа analytics вычитывает историю и фиксирует оффсеты, группа billing делает то же самое независимо, повторное подключение analytics не находит ничего нового, группа audit_replay поднимает историю целиком без фиксации позиции, а два потребителя группы delivery делят между собой три партиции.

# confluent-kafka 2.15.0, Apache Kafka 4.3.0 в KRaft, прогнано на стенде 2026-08-24
"""Чтение одного лога разными группами: независимость групп, оффсеты, повторное чтение."""
import json
import time
from confluent_kafka import Consumer, TopicPartition, KafkaError

BROKER = "localhost:9092"
TOPIC = "orders"
PARTITIONS = 3
POLL_TIMEOUT = 1.0
IDLE_LIMIT = 5  # столько пустых опросов подряд после назначения партиций считаем концом лога


def base_config(group_id: str) -> dict:
    return {
        "bootstrap.servers": BROKER,
        "group.id": group_id,
        "auto.offset.reset": "earliest",  # новая группа начинает с начала лога, а не с конца
        "enable.auto.commit": False,      # оффсеты фиксируем руками, чтобы видеть момент коммита
    }


def subscribe(consumer: Consumer, group_id: str) -> dict:
    """Подписка с колбэком назначения: пока партиции не розданы, читать нечего."""
    state = {"assigned": False}

    def on_assign(_consumer, partitions):
        state["assigned"] = True
        print(f"[{group_id}] группа получила партиции {sorted(p.partition for p in partitions)}")

    consumer.subscribe([TOPIC], on_assign=on_assign)
    return state


def drain(consumer: Consumer, label: str, state: dict, timeout: float = 60.0) -> int:
    """Читает топик до конца лога и печатает, что вычитал."""
    seen, idle = [], 0
    deadline = time.time() + timeout
    while idle < IDLE_LIMIT and time.time() < deadline: msg = consumer.poll(POLL_TIMEOUT) if msg is None: if state["assigned"]: # пустой опрос до назначения партиций концом лога не считается idle += 1 continue if msg.error(): if msg.error().code() != KafkaError._PARTITION_EOF: print(f"[{label}] ОШИБКА: {msg.error()}") continue idle = 0 payload = json.loads(msg.value()) seen.append((msg.partition(), msg.offset(), payload["order"])) print(f"[{label}] прочитано событий: {len(seen)}") for partition, offset, order in seen: print(f"[{label}] партиция {partition} оффсет {offset} заказ {order}") return len(seen) def show_committed(consumer: Consumer, group_id: str) -> None:
    """Позиция группы в логе хранится на брокере и переживает остановку потребителя."""
    parts = [TopicPartition(TOPIC, p) for p in range(PARTITIONS)]
    for tp in consumer.committed(parts, timeout=10):
        value = "нет" if tp.offset < 0 else tp.offset print(f"[{group_id}] зафиксированный оффсет партиции {tp.partition}: {value}") def read_as(group_id: str, commit: bool = True) -> int:
    consumer = Consumer(base_config(group_id))
    state = subscribe(consumer, group_id)
    count = drain(consumer, group_id, state)
    if commit and count:
        consumer.commit(asynchronous=False)  # фиксируем позицию явно
    show_committed(consumer, group_id)
    consumer.close()
    return count


def two_consumers_one_group(group_id: str) -> None:
    """Внутри группы партиции делятся между потребителями, каждая достаётся ровно одному."""
    consumers = [Consumer(base_config(group_id)) for _ in range(2)]
    for consumer in consumers:
        consumer.subscribe([TOPIC])
    deadline = time.time() + 30
    while time.time() < deadline: held = [] for consumer in consumers: consumer.poll(0.5) held.append(sorted(tp.partition for tp in consumer.assignment())) if all(held) and sum(len(h) for h in held) == PARTITIONS: break for i, consumer in enumerate(consumers, start=1): parts = sorted(tp.partition for tp in consumer.assignment()) print(f"[{group_id}] потребитель {i} держит партиции {parts}") for consumer in consumers: consumer.close() def main() -> None:
    print("=== группа analytics читает лог с начала ===")
    first = read_as("analytics")

    print("=== группа billing читает тот же лог независимо ===")
    second = read_as("billing")
    print(f"обе группы вычитали одинаковое число событий: {first == second} ({first})")

    print("=== analytics подключается снова: продолжает с зафиксированного оффсета ===")
    again = read_as("analytics")
    print(f"новых событий для analytics: {again}")

    print("=== повторное чтение истории новой группой audit_replay ===")
    replay = read_as("audit_replay", commit=False)
    print(f"replay поднял из лога событий: {replay}")

    print("=== два потребителя в одной группе делят партиции ===")
    two_consumers_one_group("delivery")


if __name__ == "__main__":
    main()

Ключевые строки прогона выглядят так.

[analytics] прочитано событий: 8
[billing] прочитано событий: 8
обе группы вычитали одинаковое число событий: True (8)
[analytics] прочитано событий: 0
новых событий для analytics: 0
[audit_replay] прочитано событий: 8
replay поднял из лога событий: 8
[delivery] потребитель 1 держит партиции [2]
[delivery] потребитель 2 держит партиции [0, 1]

Две группы забрали по 8 событий каждая, потому что чтение не расходует запись. Повторное подключение analytics дало 0 новых событий, ведь позиция группы хранится на брокере. Новая группа audit_replay подняла все 8 событий заново, и это тот самый replay в чистом виде. Внутри группы delivery три партиции разделились между двумя потребителями без пересечений.

Взгляд снаружи подтверждает картину. Команда kafka-get-offsets.sh показывает конец лога по каждой партиции, а kafka-consumer-groups.sh отдаёт позиции групп и отставание.

orders:0:3
orders:1:2
orders:2:3

GROUP           TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
analytics       orders    0          3               3               0
analytics       orders    1          2               2               0
analytics       orders    2          3               3               0

Отставание нулевое, группа догнала конец лога. Группа audit_replay в списке групп присутствует, но оффсетов у неё нет, поскольку она читала историю и ничего не фиксировала. Замечу одну грабельку из прогона. Если начать читать до того, как координатор назначил группе партиции, первое событие теряется, потому что опрос успевает вернуть его раньше назначения. Лечится это колбэком назначения и правилом не считать пустой опрос концом лога до тех пор, пока партиции не розданы. Разобраться с настройкой продюсеров и потребителей на боевых объёмах поможет курс Разработчик Apache Kafka.

 

Заключение

Потоковая передача событий отличается от очереди сообщений одним архитектурным решением, а именно тем, что событие остаётся в логе после прочтения. Отсюда растёт всё остальное, от независимых групп потребителей до повторного чтения истории новым сервисом. Порядок гарантируется в пределах партиции, поэтому выбор ключа события определяет корректность обработки не меньше, чем код обработчика. Платить за эти свойства приходится эксплуатацией кластера и дисциплиной вокруг схем данных, и на паре интеграций такой ценник не окупается. Зато когда факт события нужен пяти системам сразу и его придётся перечитывать, альтернативы у модели практически нет.

 

Референсные ссылки