Apache Kafka -- распределённая платформа потоковой передачи данных. Kafka работает как распределённый коммит-лог: записывает события в упорядоченный, неизменяемый журнал.
Основные свойства
Высокая пропускная способность -- миллионы сообщений в секунду
Долговременное хранение -- данные хранятся на диске, не только в памяти
Горизонтальная масштабируемость -- добавление брокеров без простоя
Каждая партиция имеет одного лидера (leader) и N реплик (replicas)
Партиция -- единица параллелизма
Распределение по партициям
<?php
declare(strict_types=1);
/**
* Message partitioning strategies
*/
final class PartitioningExample
{
/**
* Messages with the same key always go to the same partition.
* This guarantees ordering for related events.
*/
public function determinePartition(string $key, int $partitionCount): int
{
// Default Kafka partitioner: murmur2 hash
return abs(crc32($key)) % $partitionCount;
}
/**
* Example: all events for user_123 go to partition 2
* This means all events for this user are strictly ordered
*/
public function example(): void
{
$userId = 'user_123';
$partitions = 6;
$target = $this->determinePartition($userId, $partitions);
// user_123 -> always partition 2 (deterministic)
}
}
package kafka
import "hash/crc32"
// DeterminePartition calculates target partition for a given key.
// Messages with the same key always go to the same partition,
// guaranteeing ordering for related events.
func DeterminePartition(key string, partitionCount int) int {
hash := crc32.ChecksumIEEE([]byte(key))
return int(hash) % partitionCount
}
// Example: all events for user_123 go to a deterministic partition.
func Example() {
userID := "user_123"
partitions := 6
target := DeterminePartition(userID, partitions)
_ = target // user_123 -> always same partition (deterministic)
}
using System.IO.Hashing;
using System.Text;
// Message partitioning strategies.
public static class Partitioning
{
// DeterminePartition calculates target partition for a given key.
// Messages with the same key always go to the same partition,
// guaranteeing ordering for related events.
public static int DeterminePartition(string key, int partitionCount)
{
var hash = Crc32.HashToUInt32(Encoding.UTF8.GetBytes(key));
return (int)(hash % (uint)partitionCount);
}
// Example: all events for user_123 go to a deterministic partition.
public static void Example()
{
const string userId = "user_123";
const int partitions = 6;
var target = DeterminePartition(userId, partitions);
_ = target; // user_123 -> always same partition (deterministic)
}
}
from zlib import crc32
def determine_partition(key: str, partition_count: int) -> int:
"""Calculate target partition for a given key.
Messages with the same key always go to the same partition,
guaranteeing ordering for related events.
"""
return crc32(key.encode()) % partition_count
def example() -> None:
"""All events for user_123 go to a deterministic partition."""
user_id = "user_123"
partitions = 6
target = determine_partition(user_id, partitions)
_ = target # user_123 -> always same partition (deterministic)
> **Правило:** сообщения с одинаковым ключом попадают в одну партицию. Используйте ключ для группировки связанных событий (все заказы одного пользователя, все события одной сессии).
Consumer Groups
Consumer group -- группа потребителей, которые совместно читают топик. Каждая партиция назначается ровно одному потребителю в группе.
Topic: orders (4 partitions)
Consumer Group "order-service":
Consumer A: reads Part 0, Part 1
Consumer B: reads Part 2, Part 3
Consumer Group "analytics":
Consumer C: reads Part 0, Part 1, Part 2, Part 3
Правила назначения
Потребителей
Партиций
Результат
2
4
По 2 партиции на потребителя
4
4
По 1 партиции на потребителя
6
4
4 активных, 2 простаивают
Важно: количество потребителей в группе не должно превышать количество партиций. Лишние потребители будут простаивать.
Offsets
Offset -- позиция сообщения в партиции. Каждая consumer group хранит свой текущий offset для каждой партиции.
using System.Text;
using System.Text.Json;
using Confluent.Kafka;
// Producer wraps Kafka producer with idempotent delivery.
public sealed class KafkaEventProducer : IDisposable
{
private readonly IProducer<string, string> _producer;
private readonly string _topic;
public KafkaEventProducer(string brokers, string topic)
{
var config = new ProducerConfig
{
BootstrapServers = brokers,
Acks = Acks.All, // Wait for all replicas
MessageSendMaxRetries = 3,
RetryBackoffMs = 100,
EnableIdempotence = true, // Exactly-once producer
};
_producer = new ProducerBuilder<string, string>(config).Build();
_topic = topic;
}
// SendAsync publishes an event to Kafka with the given partition key.
// Unlike the fire-and-forget PHP/Go calls, the delivery report is awaited here.
public async Task SendAsync(string key, IDictionary<string, object> payload, CancellationToken ct = default)
{
var data = JsonSerializer.Serialize(payload);
var eventType = payload.TryGetValue("type", out var t) ? t?.ToString() : null;
var message = new Message<string, string>
{
Key = key,
Value = data,
Headers =
[
new Header("event_type", Encoding.UTF8.GetBytes(eventType ?? "unknown")),
new Header("timestamp", Encoding.UTF8.GetBytes(DateTimeOffset.UtcNow.ToUnixTimeSeconds().ToString())),
new Header("source", Encoding.UTF8.GetBytes("order-service")),
],
};
// Partition is auto-assigned from the key when none is specified
var result = await _producer.ProduceAsync(_topic, message, ct);
if (result.Status == PersistenceStatus.NotPersisted)
{
throw new InvalidOperationException($"Message delivery failed for key {key}");
}
}
// Flush waits for all outstanding messages to be delivered.
public void Flush(TimeSpan timeout) => _producer.Flush(timeout);
public void Dispose() => _producer.Dispose();
}
import json
import time
from typing import Any
from confluent_kafka import KafkaException, Producer
class KafkaEventProducer:
"""Kafka producer with idempotent delivery."""
def __init__(self, brokers: str, topic: str) -> None:
self._producer = Producer(
{
"bootstrap.servers": brokers,
"acks": "all", # Wait for all replicas
"retries": 3,
"retry.backoff.ms": 100,
"enable.idempotence": True, # Exactly-once producer
}
)
self._topic = topic
def send(self, key: str, payload: dict[str, Any]) -> None:
"""Send event to Kafka; key selects the partition."""
message = json.dumps(payload)
self._producer.produce(
topic=self._topic,
key=key,
value=message,
headers=[
("event_type", str(payload.get("type", "unknown"))),
("timestamp", str(int(time.time()))),
("source", "order-service"),
],
on_delivery=self._on_delivery,
)
# Trigger delivery callbacks
self._producer.poll(0)
@staticmethod
def _on_delivery(err: KafkaException | None, msg: Any) -> None:
if err is not None:
raise RuntimeError(f"Message delivery failed: {err}")
def flush(self, timeout: float = 10.0) -> None:
"""Flush remaining messages before shutdown."""
remaining = self._producer.flush(timeout)
if remaining > 0:
raise RuntimeError(f"Failed to flush producer: {remaining} messages pending")
### Consumer: чтение сообщений
<?php
declare(strict_types=1);
/**
* Kafka consumer with manual offset commit
*/
final class KafkaConsumer
{
private \RdKafka\KafkaConsumer $consumer;
public function __construct(
string $brokers,
string $groupId,
array $topics,
) {
$config = new \RdKafka\Conf();
$config->set('metadata.broker.list', $brokers);
$config->set('group.id', $groupId);
$config->set('auto.offset.reset', 'earliest'); // Start from beginning if no offset
$config->set('enable.auto.commit', 'false'); // Manual commit
$config->set('max.poll.interval.ms', '300000');
$config->set('session.timeout.ms', '30000');
// Rebalance callback
$config->setRebalanceCb(
function (\RdKafka\KafkaConsumer $consumer, int $err, array $partitions = null): void {
match ($err) {
RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS => $consumer->assign($partitions),
RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS => $consumer->assign(null),
default => throw new \RuntimeException(rd_kafka_err2str($err)),
};
}
);
$this->consumer = new \RdKafka\KafkaConsumer($config);
$this->consumer->subscribe($topics);
}
/**
* Main consume loop
*/
public function run(callable $handler): void
{
while (true) {
$message = $this->consumer->consume(1000); // 1 second timeout
match ($message->err) {
RD_KAFKA_RESP_ERR_NO_ERROR => $this->processMessage($message, $handler),
RD_KAFKA_RESP_ERR__PARTITION_EOF => null, // End of partition, wait
RD_KAFKA_RESP_ERR__TIMED_OUT => null, // No message, wait
default => throw new \RuntimeException(
"Kafka error: " . rd_kafka_err2str($message->err)
),
};
}
}
private function processMessage(\RdKafka\Message $message, callable $handler): void
{
try {
$payload = json_decode($message->payload, true, 512, JSON_THROW_ON_ERROR);
$handler($payload, $message->key, $message->headers);
// Commit offset only after successful processing
$this->consumer->commit($message);
} catch (\Throwable $e) {
// Log error, send to DLQ, or retry
error_log("Failed to process message: " . $e->getMessage());
}
}
}
// Usage
$consumer = new KafkaConsumer(
brokers: 'kafka1:9092,kafka2:9092',
groupId: 'order-processor',
topics: ['orders'],
);
$consumer->run(function (array $payload, string $key, ?array $headers): void {
match ($payload['type']) {
'order.created' => handleOrderCreated($payload),
'order.paid' => handleOrderPaid($payload),
'order.cancelled' => handleOrderCancelled($payload),
default => null, // Skip unknown events
};
});
package kafka
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/confluentinc/confluent-kafka-go/v2/kafka"
)
// MessageHandler processes a decoded Kafka message.
type MessageHandler func(ctx context.Context, payload map[string]any, key string) error
// Consumer reads messages from Kafka with manual offset commit.
type Consumer struct {
consumer *kafka.Consumer
}
func NewConsumer(brokers, groupID string, topics []string) (*Consumer, error) {
c, err := kafka.NewConsumer(&kafka.ConfigMap{
"bootstrap.servers": brokers,
"group.id": groupID,
"auto.offset.reset": "earliest",
"enable.auto.commit": false,
"max.poll.interval.ms": 300000,
"session.timeout.ms": 30000,
})
if err != nil {
return nil, fmt.Errorf("create consumer: %w", err)
}
if err := c.SubscribeTopics(topics, nil); err != nil {
return nil, fmt.Errorf("subscribe topics: %w", err)
}
return &Consumer{consumer: c}, nil
}
// Run starts the main consume loop, calling handler for each message.
func (c *Consumer) Run(ctx context.Context, handler MessageHandler) error {
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
msg, err := c.consumer.ReadMessage(-1)
if err != nil {
log.Printf("Consumer error: %v", err)
continue
}
var payload map[string]any
if err := json.Unmarshal(msg.Value, &payload); err != nil {
log.Printf("Failed to unmarshal message: %v", err)
continue
}
if err := handler(ctx, payload, string(msg.Key)); err != nil {
log.Printf("Failed to process message: %v", err)
continue
}
// Commit offset only after successful processing
if _, err := c.consumer.CommitMessage(msg); err != nil {
log.Printf("Failed to commit offset: %v", err)
}
}
}
// Close shuts down the consumer.
func (c *Consumer) Close() error {
return c.consumer.Close()
}
using System.Text.Json;
using Confluent.Kafka;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
// MessageHandler processes a decoded Kafka message.
public delegate Task MessageHandler(IDictionary<string, JsonElement> payload, string key, CancellationToken ct);
// KafkaConsumerService reads messages with manual offset commit.
// Hosted as IHostedService so the loop follows the app lifetime token.
public sealed class KafkaConsumerService(
string brokers,
string groupId,
IReadOnlyList<string> topics,
MessageHandler handler,
ILogger<KafkaConsumerService> logger) : BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var config = new ConsumerConfig
{
BootstrapServers = brokers,
GroupId = groupId,
AutoOffsetReset = AutoOffsetReset.Earliest, // Start from beginning if no offset
EnableAutoCommit = false, // Manual commit
MaxPollIntervalMs = 300_000,
SessionTimeoutMs = 30_000,
};
using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe(topics);
try
{
// Main consume loop
while (!stoppingToken.IsCancellationRequested)
{
var result = consumer.Consume(stoppingToken);
if (result?.Message is null)
{
continue;
}
try
{
var payload = JsonSerializer.Deserialize<Dictionary<string, JsonElement>>(result.Message.Value)
?? throw new JsonException("empty payload");
await handler(payload, result.Message.Key, stoppingToken);
// Commit offset only after successful processing
consumer.Commit(result);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
// Log error, send to DLQ, or retry
logger.LogError(ex, "Failed to process message at {Offset}", result.TopicPartitionOffset);
}
}
}
finally
{
consumer.Close();
}
}
}
import json
import logging
from collections.abc import Awaitable, Callable
from typing import Any
from confluent_kafka import Consumer, KafkaError
logger = logging.getLogger(__name__)
MessageHandler = Callable[[dict[str, Any], str], Awaitable[None]]
class KafkaEventConsumer:
"""Kafka consumer with manual offset commit."""
def __init__(self, brokers: str, group_id: str, topics: list[str]) -> None:
self._consumer = Consumer(
{
"bootstrap.servers": brokers,
"group.id": group_id,
"auto.offset.reset": "earliest", # Start from beginning if no offset
"enable.auto.commit": False, # Manual commit
"max.poll.interval.ms": 300000,
"session.timeout.ms": 30000,
}
)
self._consumer.subscribe(topics)
async def run(self, handler: MessageHandler) -> None:
"""Main consume loop."""
try:
while True:
msg = self._consumer.poll(1.0) # 1 second timeout
if msg is None:
continue # No message, wait
err = msg.error()
if err is not None:
if err.code() == KafkaError._PARTITION_EOF:
continue # End of partition, wait
raise RuntimeError(f"Kafka error: {err}")
try:
payload = json.loads(msg.value())
key = msg.key().decode() if msg.key() else ""
await handler(payload, key)
# Commit offset only after successful processing
self._consumer.commit(msg, asynchronous=False)
except Exception:
# Log error, send to DLQ, or retry
logger.exception("Failed to process message at offset %s", msg.offset())
finally:
self._consumer.close()
## Паттерны использования Kafka
Dead Letter Queue (DLQ)
Сообщения, которые не удалось обработать, отправляются в отдельный топик:
<?php
declare(strict_types=1);
final class ResilientConsumer
{
private const MAX_RETRIES = 3;
public function __construct(
private readonly KafkaConsumer $consumer,
private readonly KafkaProducer $dlqProducer,
) {}
public function processWithRetry(array $payload, string $key): void
{
$retries = (int) ($payload['_retry_count'] ?? 0);
try {
$this->processOrder($payload);
} catch (\Throwable $e) {
if ($retries >= self::MAX_RETRIES) {
// Send to Dead Letter Queue
$this->dlqProducer->send($key, [
...$payload,
'_dlq_reason' => $e->getMessage(),
'_dlq_timestamp' => time(),
'_original_topic' => 'orders',
]);
return;
}
// Retry with incremented counter
$this->dlqProducer->send($key, [
...$payload,
'_retry_count' => $retries + 1,
]);
}
}
private function processOrder(array $payload): void
{
// Business logic
}
}
package kafka
import (
"context"
"fmt"
"time"
)
const maxRetries = 3
// ResilientConsumer retries processing and sends to DLQ on failure.
type ResilientConsumer struct {
dlqProducer *Producer
}
func NewResilientConsumer(dlq *Producer) *ResilientConsumer {
return &ResilientConsumer{dlqProducer: dlq}
}
// ProcessWithRetry attempts to process an event, retrying up to maxRetries times.
func (c *ResilientConsumer) ProcessWithRetry(ctx context.Context, payload map[string]any, key string) error {
retries := 0
if v, ok := payload["_retry_count"].(float64); ok {
retries = int(v)
}
if err := c.processOrder(payload); err != nil {
if retries >= maxRetries {
// Send to Dead Letter Queue
payload["_dlq_reason"] = err.Error()
payload["_dlq_timestamp"] = time.Now().Unix()
payload["_original_topic"] = "orders"
return c.dlqProducer.Send(ctx, key, payload)
}
// Retry with incremented counter
payload["_retry_count"] = float64(retries + 1)
return c.dlqProducer.Send(ctx, key, payload)
}
return nil
}
func (c *ResilientConsumer) processOrder(payload map[string]any) error {
// Business logic
return nil
}
// ResilientConsumer retries processing and sends to DLQ on failure.
public sealed class ResilientConsumer(KafkaEventProducer dlqProducer)
{
private const int MaxRetries = 3;
public async Task ProcessWithRetryAsync(
IDictionary<string, object> payload,
string key,
CancellationToken ct = default)
{
var retries = payload.TryGetValue("_retry_count", out var raw) ? Convert.ToInt32(raw) : 0;
try
{
await ProcessOrderAsync(payload, ct);
}
catch (Exception ex)
{
if (retries >= MaxRetries)
{
// Send to Dead Letter Queue
payload["_dlq_reason"] = ex.Message;
payload["_dlq_timestamp"] = DateTimeOffset.UtcNow.ToUnixTimeSeconds();
payload["_original_topic"] = "orders";
await dlqProducer.SendAsync(key, payload, ct);
return;
}
// Retry with incremented counter
payload["_retry_count"] = retries + 1;
await dlqProducer.SendAsync(key, payload, ct);
}
}
private static Task ProcessOrderAsync(IDictionary<string, object> payload, CancellationToken ct)
{
// Business logic
return Task.CompletedTask;
}
}
import time
from typing import Any
MAX_RETRIES = 3
class ResilientConsumer:
"""Retries processing and sends to DLQ on failure."""
def __init__(self, dlq_producer: KafkaEventProducer) -> None:
self._dlq_producer = dlq_producer
def process_with_retry(self, payload: dict[str, Any], key: str) -> None:
retries = int(payload.get("_retry_count", 0))
try:
self._process_order(payload)
except Exception as exc:
if retries >= MAX_RETRIES:
# Send to Dead Letter Queue
self._dlq_producer.send(
key,
{
**payload,
"_dlq_reason": str(exc),
"_dlq_timestamp": int(time.time()),
"_original_topic": "orders",
},
)
return
# Retry with incremented counter
self._dlq_producer.send(key, {**payload, "_retry_count": retries + 1})
def _process_order(self, payload: dict[str, Any]) -> None:
# Business logic
...
### Transactional Outbox
Для гарантии согласованности между БД и Kafka:
<?php
declare(strict_types=1);
/**
* Outbox pattern: save event to DB, then relay to Kafka
*/
final class OutboxPublisher
{
public function __construct(
private readonly \PDO $db,
) {}
/**
* Save order and event in single transaction
*/
public function createOrder(array $orderData): string
{
$this->db->beginTransaction();
try {
$orderId = $this->insertOrder($orderData);
// Save event to outbox table (same transaction)
$stmt = $this->db->prepare(<<<SQL
INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES (:id, 'order', :aggregate_id, 'order.created', :payload, NOW())
SQL);
$stmt->execute([
'id' => uuid_create(),
'aggregate_id' => $orderId,
'payload' => json_encode([
'order_id' => $orderId,
'user_id' => $orderData['user_id'],
'total' => $orderData['total'],
]),
]);
$this->db->commit();
return $orderId;
} catch (\Throwable $e) {
$this->db->rollBack();
throw $e;
}
}
private function insertOrder(array $data): string
{
// Insert order into orders table
return 'ord_' . bin2hex(random_bytes(8));
}
}
package outbox
import (
"context"
"crypto/rand"
"database/sql"
"encoding/hex"
"encoding/json"
"fmt"
"github.com/google/uuid"
)
// OutboxPublisher saves events to an outbox table within the same DB transaction.
type OutboxPublisher struct {
db *sql.DB
}
func NewOutboxPublisher(db *sql.DB) *OutboxPublisher {
return &OutboxPublisher{db: db}
}
// CreateOrder saves order and outbox event in a single transaction.
func (p *OutboxPublisher) CreateOrder(ctx context.Context, orderData map[string]any) (string, error) {
tx, err := p.db.BeginTx(ctx, nil)
if err != nil {
return "", fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
orderID := generateOrderID()
// Insert order
if _, err := tx.ExecContext(ctx,
`INSERT INTO orders (id, user_id, total, status) VALUES ($1, $2, $3, 'pending')`,
orderID, orderData["user_id"], orderData["total"],
); err != nil {
return "", fmt.Errorf("insert order: %w", err)
}
// Save event to outbox table (same transaction)
payload, _ := json.Marshal(map[string]any{
"order_id": orderID,
"user_id": orderData["user_id"],
"total": orderData["total"],
})
if _, err := tx.ExecContext(ctx,
`INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES ($1, 'order', $2, 'order.created', $3, NOW())`,
uuid.New().String(), orderID, string(payload),
); err != nil {
return "", fmt.Errorf("insert outbox event: %w", err)
}
if err := tx.Commit(); err != nil {
return "", fmt.Errorf("commit tx: %w", err)
}
return orderID, nil
}
func generateOrderID() string {
b := make([]byte, 8)
rand.Read(b)
return "ord_" + hex.EncodeToString(b)
}
using System.Security.Cryptography;
using System.Text.Json;
using Npgsql;
// Outbox pattern: save event to DB, then relay to Kafka.
public sealed class OutboxPublisher(NpgsqlDataSource dataSource)
{
// CreateOrderAsync saves order and outbox event in a single transaction.
public async Task<string> CreateOrderAsync(OrderData orderData, CancellationToken ct = default)
{
await using var conn = await dataSource.OpenConnectionAsync(ct);
await using var tx = await conn.BeginTransactionAsync(ct);
var orderId = GenerateOrderId();
await using (var insertOrder = new NpgsqlCommand(
"INSERT INTO orders (id, user_id, total, status) VALUES ($1, $2, $3, 'pending')", conn, tx))
{
insertOrder.Parameters.AddWithValue(orderId);
insertOrder.Parameters.AddWithValue(orderData.UserId);
insertOrder.Parameters.AddWithValue(orderData.Total);
await insertOrder.ExecuteNonQueryAsync(ct);
}
// Save event to outbox table (same transaction)
var payload = JsonSerializer.Serialize(new
{
order_id = orderId,
user_id = orderData.UserId,
total = orderData.Total,
});
await using (var insertEvent = new NpgsqlCommand(
"""
INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES ($1, 'order', $2, 'order.created', $3::jsonb, NOW())
""", conn, tx))
{
insertEvent.Parameters.AddWithValue(Guid.NewGuid());
insertEvent.Parameters.AddWithValue(orderId);
insertEvent.Parameters.AddWithValue(payload);
await insertEvent.ExecuteNonQueryAsync(ct);
}
await tx.CommitAsync(ct);
return orderId;
}
private static string GenerateOrderId() =>
"ord_" + Convert.ToHexString(RandomNumberGenerator.GetBytes(8)).ToLowerInvariant();
}
public sealed record OrderData(string UserId, decimal Total);
import json
import secrets
import uuid
from dataclasses import dataclass
import asyncpg
@dataclass(frozen=True, slots=True)
class OrderData:
user_id: str
total: float
class OutboxPublisher:
"""Outbox pattern: save event to DB, then relay to Kafka."""
def __init__(self, pool: asyncpg.Pool) -> None:
self._pool = pool
async def create_order(self, order_data: OrderData) -> str:
"""Save order and event in single transaction."""
order_id = self._generate_order_id()
async with self._pool.acquire() as conn, conn.transaction():
await conn.execute(
"INSERT INTO orders (id, user_id, total, status) VALUES ($1, $2, $3, 'pending')",
order_id,
order_data.user_id,
order_data.total,
)
# Save event to outbox table (same transaction)
payload = json.dumps(
{
"order_id": order_id,
"user_id": order_data.user_id,
"total": order_data.total,
}
)
await conn.execute(
"""
INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES ($1, 'order', $2, 'order.created', $3::jsonb, NOW())
""",
uuid.uuid4(),
order_id,
payload,
)
return order_id
@staticmethod
def _generate_order_id() -> str:
return f"ord_{secrets.token_hex(8)}"
## Мониторинг Kafka
Ключевые метрики
Метрика
Описание
Порог алерта
Consumer Lag
Отставание потребителя
> 10000 сообщений
Under-replicated Partitions
Недореплицированные партиции
> 0
Request Latency
Задержка запросов
p99 > 100ms
Disk Usage
Использование диска
> 80%
ISR Shrink Rate
Частота выпадения из ISR
> 0/min
Итоги
Kafka -- распределённый коммит-лог с высокой пропускной способностью и отказоустойчивостью
Партиции обеспечивают параллелизм и порядок (внутри одной партиции)
Consumer groups позволяют масштабировать потребление
php-rdkafka предоставляет полноценный доступ к Kafka из PHP
Используйте manual commit, DLQ и outbox pattern для надёжной обработки