{
    "version": "https:\/\/jsonfeed.org\/version\/1.1",
    "title": "i-am-roma 🪬 analyst: заметки с тегом брокеры",
    "_rss_description": "i-am-roma 🪬 analyst",
    "_rss_language": "ru",
    "_itunes_email": "",
    "_itunes_categories_xml": "",
    "_itunes_image": "",
    "_itunes_explicit": "",
    "home_page_url": "https:\/\/i-am-roma.com\/tags\/brokery\/",
    "feed_url": "https:\/\/i-am-roma.com\/tags\/brokery\/json\/",
    "icon": false,
    "authors": [
        {
            "name": "i-am-roma",
            "url": "https:\/\/i-am-roma.com\/",
            "avatar": false
        }
    ],
    "items": [
        {
            "id": "3",
            "url": "https:\/\/i-am-roma.com\/all\/kafka-eto-ne-ochered-eto-zhurnal-kotory-nikto-ne-stiraet-2\/",
            "title": "Kafka — это не очередь. Это журнал, который никто не сотрёт",
            "content_html": "<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/cover.jpg\" width=\"2400\" height=\"1260\" alt=\"\" \/>\n<\/div>\n<p class=\"lead\">Если вы работали с очередями вроде RabbitMQ или SQS, у Kafka есть одна особенность, которая ломает привычную интуицию: прочитанное сообщение никуда не девается. Разберёмся, как это устроено и что из этого следует — на простых примерах, без воды.<\/p>\n<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/diagram-queue-vs-log.jpg\" width=\"2200\" height=\"800\" alt=\"\" \/>\n<div class=\"e2-text-caption\">Слева — как работает очередь. Справа — как работает Kafka<\/div>\n<\/div>\n<h2>Почему это не очередь<\/h2>\n<p>В очереди сообщение живёт до первого прочтения: воркер забрал — брокер удалил. Удобно для задач «сделай что-то один раз», но плохо подходит, если один и тот же поток данных нужен нескольким системам сразу.<\/p>\n<p>Kafka устроена иначе. Сообщение дописывается в конец файла на диске — такой файл называют <b>логом<\/b>, или commit log — и хранится там столько, сколько настроено, независимо от того, кто и сколько раз его уже прочитал.<\/p>\n<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/diagram-commit-log.jpg\" width=\"2200\" height=\"600\" alt=\"\" \/>\n<div class=\"e2-text-caption\">Продюсер только дописывает справа. Потребители читают с любой позиции (offset) и не удаляют прочитанное<\/div>\n<\/div>\n<p>Отсюда главный практический эффект: пять разных систем — аналитика, поиск, антифрод, уведомления, аудит — могут читать один и тот же поток независимо, каждая в своём темпе, не мешая друг другу.<br \/>\n▫️<b>O(1)<\/b> — чтение и запись это последовательный доступ к диску, не поиск по индексу<br \/>\n▫️<b>N×<\/b> — сколько угодно независимых читателей без доп. нагрузки на источник<br \/>\n▫️<b>TTL \/ ∞<\/b> — хранение настраивается: от часов до бессрочного<\/p>\n<h2>Из чего состоит кластер<\/h2>\n<p>Три термина, без которых дальше не разобраться:<\/p>\n<ul>\n<li><b>Topic<\/b> — логический поток событий, например <b>dialog-events<\/b>.<\/li>\n<li><b>Partition<\/b> — топик физически режется на партиции: независимые упорядоченные логи. Это единица параллелизма.<\/li>\n<li><b>Broker<\/b> — сервер кластера, хранящий партиции. У каждой партиции один лидер и несколько реплик-фолловеров на других брокерах.<\/li>\n<\/ul>\n<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/diagram-cluster-anatomy.jpg\" width=\"2000\" height=\"680\" alt=\"\" \/>\n<div class=\"e2-text-caption\">Топик на 3 партиции с replication factor = 2 — у каждой партиции один лидер и реплики на других брокерах<\/div>\n<\/div>\n<p>Тут есть два правила, которые часто путают:<\/p>\n<p><b>Порядок гарантирован только внутри партиции.<\/b> Между партициями — нет. Нужен строгий порядок между конкретными событиями (например, все события одного диалога) — они должны попадать в одну партицию.<\/p>\n<p><b>За отказоустойчивость отвечают реплики, а не партиции.<\/b> Продюсеры и консьюмеры общаются только с лидером. Фолловеры молча копируют его лог. Упал лидер — один из синхронных фолловеров (ISR) становится новым лидером автоматически.<\/p>\n<h2>Как не потерять порядок<\/h2>\n<p>У каждого сообщения есть необязательный <b>key<\/b>. Если он задан, Kafka всегда кладёт сообщения с одинаковым ключом в одну и ту же партицию — это единственный способ гарантировать порядок между ними.<\/p>\n<ul>\n<li>Нужен порядок внутри диалога или у одного юзера → задайте <b>key = dialog_id<\/b> (или <b>user_id<\/b>).<\/li>\n<li>Порядок не важен, важна равномерная нагрузка → не задавайте key — Kafka раскидает сообщения round-robin.<\/li>\n<\/ul>\n<p class=\"loud\">Число партиций топика нельзя безболезненно уменьшить, а увеличение задним числом рвёт гарантию порядка для уже написанных данных. Планируйте количество партиций с запасом, а не наращивайте реактивно.<\/p>\n<h2>Кто и как читает<\/h2>\n<p>Консьюмеры объединяются в <b>consumer group<\/b> по общему <b>group.id<\/b>. Правило простое: одну партицию внутри группы читает только один консьюмер одновременно — это и есть встроенный механизм масштабирования обработки.<\/p>\n<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/diagram-consumer-groups.jpg\" width=\"2000\" height=\"620\" alt=\"\" \/>\n<div class=\"e2-text-caption\">Четыре партиции распределены между двумя консьюмерами одной группы; добавление третьего вызовет rebalance<\/div>\n<\/div>\n<p>Партиций больше, чем консьюмеров — лишние простаивают. Число партиций стоит планировать исходя из желаемого параллелизма обработки, а не только из throughput записи.<\/p>\n<p>Офсет, который коммитит консьюмер — это отметка «докуда я дочитал», а не команда на удаление. Тот же топик в любой момент может независимо перечитать другая consumer group.<\/p>\n<h2>Насколько надёжна доставка<\/h2>\n<p>«Exactly-once» в Kafka — это не волшебная кнопка, а комбинация настроек, за которую вы платите производительностью. По умолчанию включён средний по надёжности вариант.<\/p>\n<div class=\"e2-text-picture\">\n<img src=\"https:\/\/i-am-roma.com\/pictures\/diagram-delivery-spectrum.jpg\" width=\"2200\" height=\"560\" alt=\"\" \/>\n<div class=\"e2-text-caption\">Три уровня гарантий доставки — от самого быстрого до самого надёжного<\/div>\n<\/div>\n<p>Две настройки, которые управляют этим спектром:<\/p>\n<ul>\n<li><b>acks<\/b> — сколько реплик должны подтвердить запись: <b>0<\/b> (никого не ждать), <b>1<\/b> (только лидер), <b>all<\/b> (все синхронные реплики).<\/li>\n<li><b>enable.idempotence<\/b> — брокер присваивает каждому сообщению номер и сам отбрасывает дубли от повторной отправки при сетевых ретраях.<\/li>\n<\/ul>\n<h2>Retention: два способа не хранить всё вечно<\/h2>\n<p><b>По умолчанию<\/b> (cleanup.policy=delete) сообщения стираются по возрасту или объёму — обычный поток событий, где старое со временем не нужно.<\/p>\n<p><b>Log compaction<\/b> (cleanup.policy=compact) — Kafka хранит только последнее значение для каждого ключа. Топик превращается в снэпшот текущего состояния, а не в историю: удобно для восстановления состояния сервиса после падения (event sourcing) или для «текущего профиля пользователя».<\/p>\n<h2>Масштабирование в двух словах<\/h2>\n<p>Разные оси дают разный эффект, и их легко перепутать:<\/p>\n<ul>\n<li><b>Больше партиций<\/b> → выше throughput записи и параллелизм чтения. Но больше файловых дескрипторов на брокер и дольше выборы лидера при отказе.<\/li>\n<li><b>Больше реплик (replication factor)<\/b> → надёжнее, но не быстрее. Клиенты всегда общаются только с лидером — реплики не ускоряют ни чтение, ни запись.<\/li>\n<\/ul>\n<p class=\"loud\">min.insync.replicas — недооценённая настройка: сколько реплик должно быть «в синхроне», чтобы запись с acks=all считалась успешной. При replication.factor=3 и min.insync.replicas=2 кластер переживёт отказ одного брокера без остановки записи; при отказе двух — запись остановится вместо тихой потери данных.<\/p>\n<h2>Экосистема вокруг Kafka<\/h2>\n<ul>\n<li><b>Kafka Connect<\/b> — готовые коннекторы для загрузки данных в Kafka и выгрузки из неё (например, Debezium для CDC из Postgres).<\/li>\n<li><b>Kafka Streams<\/b> — библиотека для потоковой обработки прямо поверх топиков, обычное JVM-приложение.<\/li>\n<li><b>ksqlDB<\/b> — те же потоковые трансформации, но на SQL, без кода.<\/li>\n<li><b>Schema Registry<\/b> — версионирование схем сообщений и проверка обратной совместимости для команд, которые пишут\/читают один топик.<\/li>\n<\/ul>\n<h2>Когда Kafka уместна, а когда нет<\/h2>\n<p><b>Хорошо подходит, если:<\/b><\/p>\n<ul>\n<li>нужно, чтобы несколько систем читали один поток независимо;<\/li>\n<li>важен большой throughput на обычном железе;<\/li>\n<li>нужна возможность переиграть историю (реплей) после бага в обработке.<\/li>\n<\/ul>\n<p><b>Плохо подходит, если:<\/b><\/p>\n<ul>\n<li>нужны SQL-запросы и фильтрация по произвольному полю — Kafka не хранилище, только последовательное чтение;<\/li>\n<li>сообщений мало и они очень мелкие — overhead на каждое делает лёгкие очереди вроде SQS эффективнее;<\/li>\n<li>нет ресурсов на операционную поддержку — свой кластер это ZooKeeper\/KRaft, мониторинг ISR и ребалансировка.<\/li>\n<\/ul>\n<table cellpadding=\"0\" cellspacing=\"0\" border=\"0\" class=\"e2-text-table\">\n<tr>\n<td>Инструмент<\/td>\n<td>Что сильнее у него<\/td>\n<td>Когда взять вместо Kafka<\/td>\n<\/tr>\n<tr>\n<td>Kafka<\/td>\n<td>throughput, независимые консьюмеры, реплей<\/td>\n<td>—<\/td>\n<\/tr>\n<tr>\n<td>RabbitMQ<\/td>\n<td>гибкая маршрутизация, приоритеты<\/td>\n<td>Task-очередь: сообщение нужно ровно одному воркеру<\/td>\n<\/tr>\n<tr>\n<td>Amazon SQS<\/td>\n<td>zero-ops, простота в AWS<\/td>\n<td>Простой pipeline без своей инфраструктуры<\/td>\n<\/tr>\n<tr>\n<td>Redis Streams<\/td>\n<td>минимальная latency<\/td>\n<td>Небольшой объём, лёгкий кейс, Redis уже в стеке<\/td>\n<\/tr>\n<tr>\n<td>Apache Pulsar<\/td>\n<td>мультиарендность<\/td>\n<td>Нужно масштабировать storage отдельно от compute<\/td>\n<\/tr>\n<\/table>\n<h2>Где Kafka встречается в реальных проектах<\/h2>\n<ol start=\"1\">\n<li><b>Шина событий между микросервисами<\/b> — вместо синхронных вызовов друг друга.<\/li>\n<li><b>CDC<\/b> — Debezium стримит изменения из Postgres\/MySQL в реальном времени.<\/li>\n<li><b>Приём потока в аналитику<\/b> — события пишутся в Kafka, оттуда попадают в ClickHouse и подобные хранилища.<\/li>\n<li><b>Log aggregation<\/b> — центральный буфер логов и метрик перед системами мониторинга.<\/li>\n<\/ol>\n<h2>Минимальный рабочий пример<\/h2>\n<pre class=\"e2-text-code\"><code class=\"bash\">kafka-topics.sh --create --topic dialog-events \\\n  --partitions 6 --replication-factor 3 \\\n  --config retention.ms=2592000000   # 30 дней\n\n# producer.properties — надёжная запись\nacks=all\nenable.idempotence=true\nlinger.ms=5   # небольшая задержка ради батчинга — выше throughput<\/code><\/pre><pre class=\"e2-text-code\"><code class=\"java\">\/\/ пишем с ключом dialog_id — порядок гарантирован внутри диалога\nproducer.send(new ProducerRecord&lt;&gt;(&quot;dialog-events&quot;, dialogId, eventJson));\n\n\/\/ коммитим offset после обработки, а не до — at-least-once\nwhile (true) {\n    var records = consumer.poll(Duration.ofMillis(500));\n    for (var record : records) process(record);\n    consumer.commitSync();\n}<\/code><\/pre>",
            "summary": "В очереди сообщение живёт до первого прочтения: воркер забрал — брокер удалил. Удобно для задач «сделай что-то один раз», но плохо подходит, если один и тот же поток данных нужен нескольким системам сразу",
            "date_published": "2026-08-11T19:16:35+04:00",
            "date_modified": "2026-08-11T20:35:51+04:00",
            "tags": [
                "kafka",
                "брокеры"
            ],
            "image": "https:\/\/i-am-roma.com\/pictures\/cover.jpg",
            "_date_published_rfc2822": "Tue, 11 Aug 2026 19:16:35 +0400",
            "_rss_guid_is_permalink": "false",
            "_rss_guid": "3",
            "_e2_data": {
                "is_favourite": false,
                "links_required": [
                    "highlight\/highlight.js",
                    "highlight\/highlight.css"
                ],
                "og_images": [
                    "https:\/\/i-am-roma.com\/pictures\/cover.jpg",
                    "https:\/\/i-am-roma.com\/pictures\/diagram-queue-vs-log.jpg",
                    "https:\/\/i-am-roma.com\/pictures\/diagram-commit-log.jpg",
                    "https:\/\/i-am-roma.com\/pictures\/diagram-cluster-anatomy.jpg",
                    "https:\/\/i-am-roma.com\/pictures\/diagram-consumer-groups.jpg",
                    "https:\/\/i-am-roma.com\/pictures\/diagram-delivery-spectrum.jpg"
                ]
            }
        }
    ],
    "_e2_version": 4199,
    "_e2_ua_string": "Aegea 11.5 (v4199)"
}