Apache Kafka Vocabulary: 30 Terms for Event Streaming Developers (англійською)

Теми, розділи, групи користувачів, зміщення, потоки Kafka, реєстр схем і словник потоків подій.

Якщо ви працюєте над серверними системами і ваша команда використовує архітектуру, керовану подією, можливо, ви вже чули про Kafka від колег. Apache Kafka стала хребетом конвеєрів даних реального часу в компаніях будь-якого розміру, і розуміння словника навколо нього є важливим — не тільки для написання коду, але і для впевненої участі в обговореннях дизайну, переглядів коду і викликів інциденту. У цій статті описано 30 основних термінів, з якими ви зіткнетеся під час роботи з Kafka і Kafka Streams, а також наведено визначення простими англійськими словами і приклади реальних розмов, які допоможуть вам зрозуміти, як їх використовувати.


Основні ринки та склади

** тема ** — канал з назвою, на який виробники надсилають повідомлення, а споживачі читають їх. Розглядайте цю категорію як категорію або подачу, подібну до теки для певного типу подій.

“Ми вирішили створити окрему ** тему ** для подій платежу, а не змішувати їх з подіями замовлення.”

«Скільки повідомлень за секунду це тема зараз обробляє?»

** розділ ** — тему розділено на один або декілька розділів, кожен з яких є упорядкованою, незмінною послідовністю записів. Розділи є первинною одиницею паралельності в Kafka.

«Ми збільшили кількість розділів з 6 до 12, щоб дозволити нам масштабувати сторону споживача»

«Повідомлення з тим же ключем завжди прибувають в той же розділ, що зберігає порядок для цього ключа.»

** replica ** — копія розділу, збереженого на іншому брокері. Репліки забезпечують відсутність помилок: якщо один з брокерів не працює, інша репліка може взяти на себе цю роботу.

«Ми запускаємо фактор реплікації 3, тому кожна реплікація знаходиться в іншій зоні доступності.»

** лідер / наслідувач ** — для кожного розділу, один з брокерів є ** лідером ** (він обробляє всі читання і записи), а інші є ** наслідувачами ** (вони реплікують дані з лідеру). Якщо лідер зазнає невдачі, послідовник обирається новим лідером.

«лідер для розділу 4 перейшов до брокера 2 після перезапуску.»

«Всі запити на виробництво і отримання йдуть до лідера — послідовники є чисто для надлишку»

** offset ** — ціле число, яке унікальним чином ідентифікує позицію запису у розділі. Споживачі стежать за своїм зсувом, щоб знати, які повідомлення вони вже обробили.

Після переривання ми вручну скидаємо offset, щоб відтворити останні дві години подій

“Завжди затверджуйте ваш ** offset ** тільки після того, як ви успішно обробили повідомлення, а не раніше.”

** lag ** — різниця між останнім зміщенням у розділі і зміщенням, до якого дійшов користувач. Висока відстань означає, що споживач відстає від виробника.

«Наша затримка піднялася до 500 тисяч повідомлень під час розгортання — нам потрібно дослідити, чому обробка сповільнилася»

«Ми створили попередження, щоб повідомити нас, якщо **затримка споживача ** перевищує 10 тисяч більше ніж на п’ять хвилин»


Промисловці, виробники, споживачі

** producer ** — програма або бібліотека, яка публікує (записує) повідомлення до теми Kafka. Виробники вирішують, на який розділ надіслати повідомлення, або за допомогою ключа, кругового обчислення, або за допомогою нетипового роздільника.

продюсер налаштований з acks=all, щоб повідомлення було записано до кожної синхронної репліки перед поверненням виклику.”

** consumer ** — програма або бібліотека, яка читає повідомлення з одного або декількох розділів тем. Споживачі опитують Kafka у власному темпі і відстежують прогрес за допомогою зсувів.

”** споживач ** є idempotent - якщо він отримує те ж саме повідомлення двічі, він не буде двічі обробляти платіж. ”

** група користувачів ** — логічне групування користувачів, які співпрацюють у читанні з теми. Кожен розділ буде призначено для одного користувача у групі, що дозволяє горизонтальне масштабування.

“Ми додали третій екземпляр до consumer group і Kafka автоматично перерозподілив розділи.”

“Якщо дві окремі служби обидві потребують читати ті ж події, дайте кожній свою власну ** групу споживачів ** - вони отримають кожне повідомлення незалежно.”

** перебалансування ** — процес, за допомогою якого Kafka перерозподіляє призначення розділів серед членів групи користувачів, що відбувається у разі приєднання, виходу або аварійного завершення роботи члена групи.

«Ми бачимо багато ** перебалансування **, тому що наші споживачі продовжують витрачати час — нам потрібно налаштувати max.poll.interval.ms. »

«Під час ** перебалансування **, споживачі призупиняють обробку, тому короткі перебалансування є критичним для конвеєрів з низькою затримкою»

** тема мертвих листів ** — тема, призначена для надсилання повідомлень, які не вдалося успішно обробити, це надає інженерам змогу перевіряти і переобробляти повідомлення без блокування головного конвеєра.

«Будь-яке повідомлення, яке не вдається десеріалізувати, направляється до **мертвої теми літера **, щоб ми могли перевірити його пізніше, не затримуючи решту потоку»

«У нас є невеликий споживач, який читає з ** dead letter topic ** і відсилає попередження до Slack, коли показники невдач є ненормальними.»


Схема і серіалізація

** Реєстр схем ** — централізоване службове вікно (зазвичай, це Реєстр конфлуєнтних схем), яке зберігає і версує схеми, використовувані для серіалізації і десеріалізації повідомлень Kafka. Він накладає правила сумісності, щоб виробники і споживачі залишалися в синхроні.

Перед об’єднанням перевірте, чи є ваша зміна схеми зворотньо сумісною з Schema Registry — ми не можемо розбити існуючі споживачі

Schema Registry відхилив нове поле, тому що воно не має типового значення і ми працюємо в режимі повної сумісності.”

Avro / Protobuf schema — два популярних формати бінарної серіалізації, що використовуються з Kafka. Avro використовує визначення схеми на основі JSON і поширений в екосистемі Confluent; Protobuf (Protocol Buffers) — бінарний формат Google, відомий своєю надійною типізацією і широкою підтримкою мов.

«Ми обрали Avro за його тісну інтеграцію з Schema Registry і його компактне двійкове кодування»

«Команда мобільних додатків віддає перевагу Protobuf, тому що він генерує сильно типізовані класи як у Swift, так і в Kotlin з тієї ж схеми»

** idmpowered producer ** — виробник, налаштований так, щоб повторні спроби не призвели до дублювання повідомлень. Kafka присвоює кожному виробнику унікальний ідентифікатор і номери послідовностей для виявлення і відкидання дублікатів на стороні брокера.

«Увімкніть параметр idempotent producer — це майже нічого не коштує і запобігає дублювання записів під час перехідних помилок мережі»

** семантику точного одного разу ** — гарантія доставки, яка забезпечує обробку кожного повідомлення точно один раз, навіть у разі помилок і повторних спроб. Досягнуто у Kafka за допомогою поєднання імемпотентних виробників, транзакцій і управління компенсацією споживача.

«Для конвеєра розрахунків нам потрібна ** саме-один раз семантика ** - дублікат платежі є неприйнятними, навіть якщо брокер збивається в середині транзакції.»

«Exactly-once має певні витрати на продуктивність, тому ми вмикаємо його тільки для фінансових подій, а не для аналітичних тем»


Kafka Streams and Stream Processing (англійською)

** Kafka Streams ** — бібліотека Java (частина проекту Apache Kafka) для створення програм обробки потоків з можливістю зміни стану, які читають і записують теми Kafka без потреби у окремому кластері обробки.

«Ми замінили нашу роботу Flink на Kafka Streams, тому що ми хотіли уникнути роботи окремого кластера — він працює всередині нашого існуючого сервісу Spring Boot»

** KStream ** — абстракція потоків Kafka, що представляє неперервний, необмежений потік записів. Кожен запис розглядається як незалежна подія.

«Ми використовуємо KStream для необроблених подій клацання, тому що кожне клацання є окремим, безстатевим фактом.»

** KTable ** — абстракція потоків Kafka, що представляє потоки журналу змін, де кожен запис є upsert (вставлення або оновлення), що має унікальний ідентифікатор. KTable моделює останній стан кожного з ключів.

«Дані профіль користувача моделюються як KTable — коли користувач оновлює свою електронну пошту, новий запис замінює старий.»

“Ми приєднуємо замовлення KStream проти клієнта KTable, щоб збагатити кожен замовлення поточним рівнем клієнта.”

** windowing ** — метод групування потоку подій у скінченночасові контейнери (вікна) для агрегування. Серед типових типів вікон можна назвати вікна, що перевертаються (фіксовані, не перекриваються), вікна, що перескакують (фіксований розмір, перекриваються) і вікна сеансів (перериваються залежно від активності).

«Ми використовуємо п’ятихвилинне вікно ** window ** для підрахунку переглядів сторінок на користувача для нашого панелі в реальному часі.»

«Session windowing group all actions within a user’s visit, closing the window after 30 minutes of inactivity.»

** store of state ** — локальне, постійне сховище ключів і значень, яке підтримується завданням Kafka Streams для підтримки операцій, що залежать від стану, таких як з’ єднання і агрегації. Сховища станів резервуються у стисненій темі журналу змін.

«Агрегація підтримується RocksDB state store — переконайтеся, що екземпляр має достатньо місця на диску під час відновлення.»

** інтерактивні запитуванні ** — функція потоків Kafka, яка надає вам змогу безпосередньо запитувати сховища станів запущеної програми, перетворюючи ваш процесор потоків на мікросервіс з можливістю запитів.

«Ми виставляємо поточні суми через ** інтерактивні запити **, тому панель може отримувати підрахунки в реальному часі без втручання в базу даних.»


Як використовувати їх у розмові

Зрозуміти ці терміни — це лише половина справи — використовувати їх природно у повсякденному спілкуванні — це те, що робить вас вільним членом інженерної команди, яка зосереджена на Kafka. Ось деякі типові сценарії:

** Сценарій 1 — Обговорення архітектури: ** Ваша команда вирішує, як розділити типи подій. Ви можете сказати: « Я рекомендую окремі ** теми ** для подій, створених за допомогою замовлення і оновлених за допомогою замовлення. Таким чином, кожна ** група споживачів ** може підписатися тільки на ті події, які її цікавлять, і ми уникаємо небажаної обробки надлишків. ”

Сценарий 2 — Расследование инцидента: Служба нижнього рівня надсилає попередження. Ви можете сказати: “Споживач ** відстає ** на тему оплати досягла 200 тисяч. Давайте перевіримо, чи є причиною перебалансування петлі — я помітив, що підрахунок подів ненадовго впав під час розгортання»

** Сценарій 3 — Перегляд коду: ** Переглядаючи налаштування виробника колеги, ви коментуєте: « Вам слід увімкнути тут прапорець idempotent producer. Без цього параметра повторна спроба після перевищення часу очікування мережі може призвести до дублювання записів. І, враховуючи, що це фінансовий потік, ми дійсно повинні розглянути ** саме один раз семантика ** end-to-end. ”

** Сценарій 4 — Введення нового члена команди: ** Пояснюючи модель даних, ви кажете: « Ми моделюємо поточний стан підписки кожного користувача як ** KTable **, і ми об’ єднуємо його з вхідною подією ** KStream **, щоб вирішити, надати або відмовити у доступі. Будь-яка подія, яку ми не можемо десеріалізувати, переходить до мертвої теми літери для вручну перегляду


Краткий справочник

TermWhat it isWhy it matters
topicNamed message channelOrganises events by type
partitionOrdered sub-division of a topicEnables parallel consumption
offsetPosition of a record in a partitionTracks consumer progress
consumer groupGroup of consumers sharing partition loadHorizontal scaling
lagGap between latest and consumed offsetKey health indicator
Schema RegistryCentralised schema versioning servicePrevents producer/consumer mismatch
KTableChangelog stream (latest value per key)Models current state
exactly-once semanticsNo duplicates, no data loss guaranteeCritical for financial pipelines
dead letter topicHolding area for unprocessable messagesPrevents pipeline blockage
rebalancingRedistribution of partitions in a groupTriggered by membership changes

Завдяки цьому словнику ви зможете більш впевнено брати участь у прийнятті рішень щодо архітектури, писати більш чіткі повідомлення про збереження і runbooks, а також точніше спілкуватися під час інциденту. Наступним кроком є вправляння у використанні цих термінів у контексті — перегляд коду, розробка документів і командні зустрічі — це чудові можливості для підтвердження того, що ви вже вивчили.

Поширені запитання

Про що ця стаття "Apache Kafka Vocabulary: 30 Terms for Event Streaming Developers (англійською)"?

Теми, розділи, групи користувачів, зміщення, потоки Kafka, реєстр схем і словник потоків подій.

Чи безкоштовна ця стаття?

Так. Усі статті на CoderSlingo, включно з цією, доступні безкоштовно без реєстрації.

Скільки часу займає читання "Apache Kafka Vocabulary: 30 Terms for Event Streaming Developers (англійською)"?

Приблизно 9 min.