Core Concepts › Kafka

Kafka

Package: packages/laravel-kafka (nestlaravel/kafka, namespace NestLaravel\Kafka), installed in every generated service. The gateway carries an equivalent copy under App\Infrastructure\Kafka / App\Messaging.

#Publishing (transactional outbox)

use NestLaravel\Kafka\Contracts\EventBus;

DB::transaction(function () use (class="tk-v">$order, class="tk-v">$events) {
    class="tk-v">$order->save();
    class="tk-v">$events->publish(new OrderCreated((string) class="tk-v">$order->id, [class="tk-s">'total' => class="tk-v">$order->total]));
});

OutboxEventBus writes the event to outbox_messages inside the caller's transaction (no dual-write problem). php artisan messaging:outbox-publish --daemon (its own container/process) ships rows to Kafka:

  1. produce() each row, then flush() — the producer callback throws on any failed delivery;
  2. only then are rows marked published; failures are rescheduled (KAFKA_OUTBOX_RETRY_DELAY) up to KAFKA_OUTBOX_MAX_ATTEMPTS, then dead-lettered (<topic>.dlq) and marked failed.

Producer defaults: acks=all, enable.idempotence=true, compression=lz4, delivery timeout 30 s. Run one outbox publisher per service. Delivery is at-least-once; consumers deduplicate.

#Consuming

php artisan kafka:consume orders.events "App\Modules\Payments\Infrastructure\Messaging\OrderCreatedHandler"

Pipeline per message: deserialize → validate envelope → schema version check (KAFKA_EVENT_MAX_VERSION) → idempotency check (event_id per topic) → handler → remember → commit offset.

SituationBehaviour
Handler throws (transient)retried KAFKA_CONSUMER_MAX_RETRIES times with exponential backoff, then → <topic>.dlq
Malformed JSON / missing fields / newer schema version (poison)straight to the DLQ, no retries
Duplicate deliveryskipped (idempotency store: the app cache — use Redis in production)
DLQ publish itself failsexception; offset is not committed → message is redelivered, never lost
SIGTERM/SIGINTcurrent message finishes, offset committed, consumer leaves the group (fast rebalance)

Offsets: enable.auto.commit=false; offsets are stored/committed only after success (at-least-once). Consumer groups default to the service name (KAFKA_GROUP_ID); scale by running more consumer processes (≤ partitions). Rebalancing is handled by librdkafka.

#Topics & events

npx nestlaravel generate kafka-topic order-events --service orders --create   # + order-events.dlq
npx nestlaravel generate kafka-event order.created --service orders --consumer

#Security & operations

SettingDevProduction
KAFKA_SECURITY_PROTOCOLplaintextsasl_ssl (or ssl with client certs)
KAFKA_SASL_*, KAFKA_SSL_*–from your secret store, one principal per service
Broker ACLs–service may WRITE only its own topics, READ only topics it consumes (+ its group); DLQ likewise
auto.create.topics.enabletrue (compose)false; provision topics with partitions/retention explicitly
Replication1≥ 3, min.insync.replicas=2

Observability: watch consumer lag (kafka-consumer-groups.sh --describe or Kafka UI: docker compose --profile tools up -d kafka-ui), DLQ topic depth, outbox pending/failed row counts, and Log::critical "Failed to publish … to DLQ".

#Local development

docker-compose.infra.yml runs a single-node KRaft broker (apache/kafka:4.0.0; no ZooKeeper) on 127.0.0.1:9092 (host) / kafka:29092 (containers). Without php-rdkafka on the host, apps use the log driver (JSONL under storage/framework/kafka) so you can develop and test without a broker; nestlaravel dev --docker runs the real thing.

Edit this page on GitHub