Kafka Event Streaming
1. Overview
BANA uses Apache Kafka for asynchronous, event-driven communication between microservices. Events are published when business state changes occur (e.g. payment succeeded, stock adjusted) and consumed by downstream services that react to those changes.
All topic names are defined centrally in @nx/core at src/queues/kafka/topics.ts.
2. Topic Registry
Topic naming: nx.bana.<evt|cmd|cdc|dlq>.<domain>[.<entity>].<event>. The application-event registry (ServiceTopicDefs in packages/core/src/common/kafka/registry.ts):
| Topic | Type | Description |
|---|---|---|
nx.bana.evt.payment.succeeded | evt | Sale payment completed |
nx.bana.evt.inventory.purchase-order.received | evt | Purchase order received into inventory |
nx.bana.evt.inventory.stock.issued-for-sale | evt | Stock issued for a sale |
nx.bana.evt.inventory.stock.adjusted | evt | Stock adjusted |
nx.bana.evt.kitchen.ticket-item.status-changed | evt | Kitchen ticket item status changed |
nx.bana.evt.inventory.material.transferred | evt | Material transferred between locations |
nx.bana.evt.inventory.material.stock-changed | evt | Material stock changed |
nx.bana.evt.signal.activity.notified | evt | Activity notification (→ WebSocket) |
nx.bana.cmd.ledger.generate | cmd | Ledger generation (ledger-internal) |
Merchant, product and product-variant changes are not evt topics - they propagate as Debezium CDC topics
nx.bana.cdc.<schema>.<Table>(see CDC). There is nocommerce.initializedtopic.
3. Producer & Consumer Matrix
| Topic | Producer | Consumer(s) |
|---|---|---|
nx.bana.evt.payment.succeeded | sale | finance, inventory, invoice |
nx.bana.evt.inventory.purchase-order.received | inventory | finance |
nx.bana.evt.inventory.stock.issued-for-sale | inventory | finance |
nx.bana.evt.inventory.stock.adjusted | inventory | finance |
nx.bana.evt.kitchen.ticket-item.status-changed | sale | inventory |
nx.bana.evt.inventory.material.transferred | inventory | inventory (cross-location) |
nx.bana.evt.inventory.material.stock-changed | inventory | (no consumer yet) |
nx.bana.evt.signal.activity.notified | (no producer yet) | signal |
nx.bana.cmd.ledger.generate | ledger | ledger (internal) |
4. Message Flow Diagram
5. Consumer Groups
| Service | Client ID | Group ID |
|---|---|---|
| Finance | SVC-00040-FINANCE_CONSUMER | SVC-00040-FINANCE_CONSUMER_GROUP |
| Inventory | SVC-00050-INVENTORY_CONSUMER | SVC-00050-INVENTORY_CONSUMER_GROUP |
| Ledger | SVC-00060-LEDGER_CONSUMER | SVC-00060-LEDGER_GROUP |
| Search (CDC) | cdc-consumer | cdc-consumer-group |
6. Consumer Configurations
6.1. Finance Consumer
| Setting | Value |
|---|---|
| Topics | nx.bana.evt.payment.succeeded, …inventory.purchase-order.received, …inventory.stock.issued-for-sale, …inventory.stock.adjusted, CDC public.Merchant |
| Auto-commit | enabled |
| Fallback Mode | earliest |
| Deserializer | String keys, JSON values |
6.2. Inventory Consumer
| Setting | Value |
|---|---|
| Topics | nx.bana.evt.payment.succeeded, …kitchen.ticket-item.status-changed, …inventory.material.transferred, CDC public.Merchant, CDC public.ProductVariant |
| Fallback Mode | earliest |
6.3. Ledger Consumer
| Setting | Value |
|---|---|
| Topics | nx.bana.cmd.ledger.generate |
| Auto-commit | false (manual commit after processing) |
| Consumer Count | Configurable via APP_ENV_KAFKA_CONSUMER_COUNT |
| Acks (Producer) | ALL |
| Idempotent (Producer) | true |
6.4. Search CDC Consumer
| Setting | Value |
|---|---|
| Topics | 27 Debezium CDC topics |
| Auto-commit | false (manual offset management) |
| Max Batch Size | 200 |
| Flush Interval | 2000ms |
| Max Wait Time | 500ms |
| Max Bytes | 5MB |
See CDC / Debezium for details.
7. Event Payload Examples
7.1. Payment Success (nx.bana.evt.payment.succeeded)
{
"saleOrderId": "987654321",
"saleOrderNumber": "SO-0001",
"merchantId": "123456789",
"saleChannelId": "555",
"inventoryLocationId": "loc-1",
"items": [{ "id": "i1", "itemType": "PRODUCT_VARIANT", "itemId": "pv-1", "quantity": 2, "mode": "DEDUCT" }]
}7.2. Ledger Generate (nx.bana.cmd.ledger.generate)
{
"ledgerId": "777888999",
"ledgerType": "S1a-HKD",
"merchantId": "444555666",
"period": "2026-Q1",
"version": 1,
"isRetry": false,
"enqueueTime": 1711785600000
}8. Best Practices
8.1. Idempotency
All consumers implement idempotency checks to safely handle duplicate messages:
| Service | Strategy |
|---|---|
| Finance | Check existing transaction by merchantId + referenceType + referenceId |
| Inventory | Check existing default location by merchantId |
| Ledger | Check existing pending job by ledgerId |
8.2. Error Handling
| Strategy | Used By |
|---|---|
| Skip + Log | Finance (logs error, continues processing) |
| DLQ (Dead Letter Queue) | Search CDC (nx.bana.dlq.cdc) |
| Retry via API | Ledger (manual retry endpoint) |
| Startup Recovery | Ledger (re-enqueues stalled jobs) |
8.3. Serialization
| Direction | Format |
|---|---|
| Producer | JSON serializer (values), string serializer (keys) |
| Consumer | JSON deserializer (values), string deserializer (keys) |
9. Infrastructure
9.1. Development (Docker Compose)
3-node Kafka cluster:
nx-kafka-1:29092nx-kafka-2:29092nx-kafka-3:29092
9.2. Environment Variables
| Variable | Description |
|---|---|
APP_ENV_KAFKA_BROKERS | Comma-separated broker addresses |
APP_ENV_KAFKA_CLIENT_ID | Service-specific client identifier |
APP_ENV_KAFKA_GROUP_ID | Consumer group identifier |
APP_ENV_KAFKA_CONSUMER_COUNT | Number of parallel consumers (Ledger) |
APP_ENV_KAFKA_SASL_MECHANISM | SASL mechanism (e.g., SCRAM-SHA-512) |
APP_ENV_KAFKA_SASL_USERNAME | SASL username |
APP_ENV_KAFKA_SASL_PASSWORD | SASL password |
10. Related Documentation
| Document | Description |
|---|---|
| CDC / Debezium | Change Data Capture via Kafka |
| Finance Service | Finance event consumers |
| Inventory Service | Inventory event consumers |
| Ledger Service | Ledger generation pipeline |
| Search Service | CDC consumer for Typesense sync |