Пошук уроків, статей та іншого контенту
Пояснить модель Kafka, партиції, офсети, групи споживачів і масштабування потокової обробки подій.
Kafka — це розподілений журнал подій. Виробники записують події, а споживачі читають їх у потрібному темпі.
Основні сутності:
топік — логічна назва потоку подій;
партиція — упорядкований журнал усередині топіка;
офсет — позиція повідомлення в конкретній партиції;
producer — компонент, який записує події;
consumer — компонент, який читає події;
consumer group — набір споживачів, що спільно обробляють топік.
Подія Kafka зазвичай містить:
ключ (key);
значення (value);
офсет;
час запису;
заголовки (headers).
Kafka не видаляє повідомлення одразу після читання. Події зберігаються відповідно до політики retention — за часом або за обсягом. Тому кілька незалежних споживачів можуть читати той самий топік зі своїми позиціями.
Топік складається з однієї або кількох партицій:
orders
├── partition 0: offset 0, 1, 2, 3, ...
├── partition 1: offset 0, 1, 2, 3, ...
└── partition 2: offset 0, 1, 2, 3, ...Партиція є append-only журналом: нові записи додаються в кінець, а кожна подія отримує монотонний офсет у межах цієї партиції.
Kafka гарантує порядок лише в межах однієї партиції.
Якщо події потрапили в різні партиції, Kafka не гарантує їхній глобальний порядок:
partition 0: A1, A3
partition 1: A2Не можна робити висновок, що A1, A2, A3 були прочитані саме в такому порядку.
Якщо для сутності важливий порядок подій, усі події цієї сутності потрібно спрямовувати в одну партицію. Зазвичай для цього використовують ключ:
key = "customer-42"Producer хешує ключ і використовує результат для вибору партиції. Тому події з однаковим ключем зазвичай потрапляють в одну партицію.
Це дає порядок для конкретного клієнта, замовлення або рахунку, але не для всього топіка.
Кількість партицій визначає потенційну паралельність читання. Якщо топік має три партиції, одночасно ефективно обробляти його в одній consumer group можуть не більше трьох активних consumers.
Однак збільшення кількості партицій має наслідки:
змінюється розподіл ключів між партиціями;
глобального порядку все одно немає;
нові події можуть маршрутизуватися інакше;
старі події не переміщуються автоматично.
Кількість партицій потрібно планувати з урахуванням:
потрібної пропускної здатності;
максимальної кількості паралельних consumers;
вимог до порядку;
очікуваного зростання системи.
Офсет — це числова позиція повідомлення в конкретній партиції.
Наприклад:
partition 0:
offset 0 -> подія A
offset 1 -> подія B
offset 2 -> подія CОфсети локальні для партиції. Офсет 10 у partition 0 не пов’язаний з офсетом 10 у partition 1.
Consumer читає події, рухаючись від одного офсету до наступного. Kafka не зберігає для кожного consumer окрему копію поточного положення в пам’яті. Натомість позиція групи зберігається як службові дані Kafka.
Важливо розрізняти:
поточну позицію consumer — де процес перебуває під час читання;
зафіксований офсет — останню позицію, яку consumer повідомив Kafka;
наступний офсет — позицію, з якої потрібно продовжити читання.
Якщо consumer обробив подію, але завершився до фіксації офсету, після перезапуску подія може бути прочитана повторно.
Якщо consumer зафіксував офсет до завершення обробки, а потім завершився, подія може більше не бути прочитана цією групою.
Consumer group — це логічний ідентифікатор споживачів, які спільно обробляють топік.
Правило розподілу:
У межах однієї consumer group конкретну партицію в певний момент часу обробляє лише один consumer.
Наприклад, є топік із трьома партиціями:
orders: P0, P1, P2І consumer group із двома consumers:
consumer-1 -> P0, P1
consumer-2 -> P2Із трьома consumers:
consumer-1 -> P0
consumer-2 -> P1
consumer-3 -> P2Якщо додати четвертий consumer, він не отримає партицію:
consumer-4 -> нічогоВін залишатиметься неактивним, доки не звільниться партиція або не зміниться склад групи.
Два consumers з однаковим group.id ділять роботу:
group: billing
consumer A -> частина подій
consumer B -> інша частина подійДва consumers із різними group.id отримують незалежні копії потоку:
group: billing -> читає всі події
group: analytics -> також читає всі подіїЦе дозволяє одному топіку одночасно живити різні підсистеми:
оплату;
аналітику;
пошук;
сповіщення.
Групи мають окремі офсети, тому одна група може відставати, а інша — читати потік майже в реальному часі.
Kafka розподіляє партиції між consumers групи. Коли склад групи змінюється, відбувається ребалансування.
Причини ребалансу:
новий consumer приєднався до групи;
consumer завершився;
consumer перестав надсилати heartbeat;
змінилася підписка на топіки;
змінилася кількість партицій.
Під час ребалансу партиції можуть тимчасово перестати оброблятися. Після нього Kafka призначає їх новим consumers.
Розподіл може виглядати так:
До:
consumer-1 -> P0, P1
consumer-2 -> P2, P3
Після завершення consumer-2:
consumer-1 -> P0, P1, P2, P3Обробник повідомлень має бути готовим до повторної доставки. Навіть якщо ребалансування працює коректно, повідомлення, оброблені без зафіксованого офсету, можуть бути прочитані ще раз.
Для тривалих операцій важливо:
не блокувати consumer надовго;
підтримувати heartbeat;
налаштовувати тайм-аути відповідно до тривалості обробки;
робити обробку ідемпотентною.
Пропускна здатність consumer group обмежена кількістю партицій:
активні паралельні обробники <= кількість партиційЯкщо топік має 12 партицій, можна запустити до 12 активних consumers у групі й отримати паралельну обробку кожної партиції окремо.
Але збільшення кількості consumers не завжди лінійно збільшує продуктивність. Обмеженнями можуть бути:
база даних;
зовнішній HTTP-сервіс;
мережа;
CPU;
дискова пропускна здатність;
нерівномірний розподіл ключів.
Якщо один ключ трапляється значно частіше за інші, відповідна партиція може стати перевантаженою:
customer-1 -> P0: 90% усіх подій
інші ключі -> P1, P2, P3: 10%Додаткові consumers не вирішать проблему, якщо гаряча партиція все одно призначена лише одному consumer.
Можливі підходи:
переглянути ключ партиціонування;
розділити надмірно великий потік на логічні підпотоки;
дозволити паралельність на рівні самої сутності, якщо порядок для неї не є обов’язковим;
масштабувати downstream-сервіс, який став вузьким місцем.
Producer може:
явно вказати номер партиції;
передати ключ;
не передавати ключ, тоді Kafka використовує доступний механізм розподілу записів між партиціями.
Рекомендація:
використовуйте стабільний ключ, якщо потрібен порядок для сутності;
не використовуйте випадковий ключ, якщо всі події однієї сутності мають оброблятися послідовно;
не використовуйте один і той самий ключ для всіх подій, якщо потрібне масштабування.
Нижче наведено приклад на Node.js з KafkaJS. Він демонструє:
запис подій у топік;
використання ключа для збереження порядку подій клієнта;
запуск consumer group;
розподіл партицій між кількома процесами.
Спочатку створіть топік із трьома партиціями. Kafka broker має бути доступним за адресою localhost:9092.
kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic orders \
--partitions 3 \
--replication-factor 1Створіть Node.js-проєкт і встановіть клієнт:
npm init -y
npm install kafkajsФайл app.js:
const { Kafka } = require("kafkajs");
const kafka = new Kafka({
clientId: "orders-example",
brokers: ["localhost:9092"],
});
const topic = "orders";
async function produce() {
const producer = kafka.producer();
await producer.connect();
const orders = [
{ id: "order-1", customerId: "customer-1", status: "created" },
{ id: "order-2", customerId: "customer-2", status: "created" },
{ id: "order-1", customerId: "customer-1", status: "paid" },
{ id: "order-3", customerId: "customer-1", status: "created" },
];
await producer.send({
topic,
messages: orders.map((order) => ({
// Однаковий customerId спрямовує події клієнта в одну партицію.
key: order.customerId,
value: JSON.stringify(order),
})),
});
await producer.disconnect();
}
async function consume(groupId) {
const consumer = kafka.consumer({ groupId });
await consumer.connect();
await consumer.subscribe({
topic,
fromBeginning: true,
});
console.log(`Consumer group "${groupId}" запущено`);
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const order = JSON.parse(message.value.toString());
console.log({
topic,
partition,
offset: message.offset,
customerId: order.customerId,
orderId: order.id,
status: order.status,
});
// Тут має бути бізнес-обробка події.
// Обробник повинен коректно працювати при повторній доставці.
},
});
}
async function main() {
const mode = process.argv[2];
if (mode === "produce") {
await produce();
return;
}
if (mode === "consume") {
const groupId = process.argv[3] || "orders-workers";
await consume(groupId);
return;
}
console.log(
"Використання: node app.js produce | consume [groupId]"
);
}
main().catch((error) => {
console.error(error);
process.exitCode = 1;
});Запис подій:
node app.js produceЗапустіть два consumers з однаковою групою в різних терміналах:
node app.js consume orders-workers
node app.js consume orders-workersВони розділять партиції між собою. Якщо запустити consumer з іншою групою:
node app.js consume analytics-workersця група отримає власний незалежний потік подій.
У прикладі customer-1 використовується як ключ. Події цього клієнта потрапляють в одну партицію, тому consumer бачить їх у порядку запису в межах цієї партиції.
Типова модель Kafka consumer — at-least-once delivery:
consumer читає подію;
виконує обробку;
фіксує офсет;
переходить до наступної події.
Якщо процес завершиться між кроками 2 і 3, подія буде оброблена повторно.
Це безпечніше за фіксацію офсету до обробки, але вимагає ідемпотентного бізнес-коду.
Операція є ідемпотентною, якщо повторне виконання не змінює кінцевий результат.
Небезпечний приклад:
отримати подію "payment-created"
завжди збільшити баланс на 100При повторній доставці баланс збільшиться двічі.
Надійніший варіант:
отримати eventId
якщо eventId уже оброблявся — пропустити
інакше:
виконати операцію
зберегти eventId як обробленийПеревірка дубліката та бізнес-операція мають бути узгоджені. Якщо вони виконуються в різних незалежних транзакціях, між ними може виникнути вікно для повторної обробки.
У KafkaJS у прикладі використовується стандартне керування офсетами. Це зручно, але не означає, що зовнішня бізнес-операція і фіксація офсету є однією атомарною транзакцією.
Consumer lag — це відставання consumer group від останньої доступної події в партиції.
Спрощено:
lag = latest offset - committed offsetLag потрібно аналізувати окремо для кожної партиції та групи.
Великий або зростаючий lag може означати:
consumer обробляє події повільніше, ніж producer їх записує;
недостатньо партицій або consumers;
одна партиція перевантажена;
downstream-сервіс працює нестабільно;
часто відбуваються ребаланси;
обробник виконує надто багато синхронних операцій.
Разовий lag після короткого піку навантаження не обов’язково є проблемою. Небезпечніший сценарій — коли lag постійно зростає.
Масштабування має сенс лише тоді, коли:
у топіку є вільні партиції;
consumers справді є вузьким місцем;
downstream-системи можуть витримати додаткову паралельність.
Партиції обробляються паралельно. Порядок гарантований лише всередині однієї партиції.
Якщо consumers більше, ніж партицій, частина процесів не матиме роботи. Для масштабування потрібно збільшувати не лише кількість процесів, а й кількість партицій, якщо це сумісно з вимогами до ключів і порядку.
Якщо події одного замовлення записуються з різними ключами, вони можуть потрапити в різні партиції. Тоді порядок подій для цього замовлення не гарантований.
Повторна доставка є нормальною частиною at-least-once обробки. Операції з платежами, балансами, залишками та іншими побічними ефектами повинні враховувати дублікати.
Це може призвести до втрати події для цієї consumer group, якщо процес завершиться після фіксації, але до завершення бізнес-операції.
Довга обробка без належного керування heartbeat може спричинити вилучення consumer із групи та ребалансування. Після цього вже оброблені, але не зафіксовані події можуть бути доставлені повторно.
Велика кількість партицій не допоможе, якщо майже весь трафік потрапляє в одну партицію через невдалий ключ.
Топік Kafka складається з однієї або кількох партицій.
Партиція — це впорядкований журнал із локальними офсетами.
Глобального порядку між партиціями немає.
Одна consumer group розподіляє партиції між своїми consumers.
У межах групи одну партицію в певний момент обробляє один consumer.
Максимальна паралельність групи обмежена кількістю партицій.
Однаковий ключ допомагає зберегти порядок подій конкретної сутності.
Consumer groups із різними ідентифікаторами читають незалежні копії потоку.
Ребалансування змінює розподіл партицій і може спричинити повторну доставку.
За at-least-once моделі обробники повинні бути ідемпотентними.
Lag показує, наскільки consumer group відстає від потоку подій.