Skip to content

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):

TopicTypeDescription
nx.bana.evt.payment.succeededevtSale payment completed
nx.bana.evt.inventory.purchase-order.receivedevtPurchase order received into inventory
nx.bana.evt.inventory.stock.issued-for-saleevtStock issued for a sale
nx.bana.evt.inventory.stock.adjustedevtStock adjusted
nx.bana.evt.kitchen.ticket-item.status-changedevtKitchen ticket item status changed
nx.bana.evt.inventory.material.transferredevtMaterial transferred between locations
nx.bana.evt.inventory.material.stock-changedevtMaterial stock changed
nx.bana.evt.signal.activity.notifiedevtActivity notification (→ WebSocket)
nx.bana.cmd.ledger.generatecmdLedger 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 no commerce.initialized topic.

3. Producer & Consumer Matrix

TopicProducerConsumer(s)
nx.bana.evt.payment.succeededsalefinance, inventory, invoice
nx.bana.evt.inventory.purchase-order.receivedinventoryfinance
nx.bana.evt.inventory.stock.issued-for-saleinventoryfinance
nx.bana.evt.inventory.stock.adjustedinventoryfinance
nx.bana.evt.kitchen.ticket-item.status-changedsaleinventory
nx.bana.evt.inventory.material.transferredinventoryinventory (cross-location)
nx.bana.evt.inventory.material.stock-changedinventory(no consumer yet)
nx.bana.evt.signal.activity.notified(no producer yet)signal
nx.bana.cmd.ledger.generateledgerledger (internal)

4. Message Flow Diagram

5. Consumer Groups

ServiceClient IDGroup ID
FinanceSVC-00040-FINANCE_CONSUMERSVC-00040-FINANCE_CONSUMER_GROUP
InventorySVC-00050-INVENTORY_CONSUMERSVC-00050-INVENTORY_CONSUMER_GROUP
LedgerSVC-00060-LEDGER_CONSUMERSVC-00060-LEDGER_GROUP
Search (CDC)cdc-consumercdc-consumer-group

6. Consumer Configurations

6.1. Finance Consumer

SettingValue
Topicsnx.bana.evt.payment.succeeded, …inventory.purchase-order.received, …inventory.stock.issued-for-sale, …inventory.stock.adjusted, CDC public.Merchant
Auto-commitenabled
Fallback Modeearliest
DeserializerString keys, JSON values

6.2. Inventory Consumer

SettingValue
Topicsnx.bana.evt.payment.succeeded, …kitchen.ticket-item.status-changed, …inventory.material.transferred, CDC public.Merchant, CDC public.ProductVariant
Fallback Modeearliest

6.3. Ledger Consumer

SettingValue
Topicsnx.bana.cmd.ledger.generate
Auto-commitfalse (manual commit after processing)
Consumer CountConfigurable via APP_ENV_KAFKA_CONSUMER_COUNT
Acks (Producer)ALL
Idempotent (Producer)true

6.4. Search CDC Consumer

SettingValue
Topics27 Debezium CDC topics
Auto-commitfalse (manual offset management)
Max Batch Size200
Flush Interval2000ms
Max Wait Time500ms
Max Bytes5MB

See CDC / Debezium for details.

7. Event Payload Examples

7.1. Payment Success (nx.bana.evt.payment.succeeded)

json
{
  "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)

json
{
  "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:

ServiceStrategy
FinanceCheck existing transaction by merchantId + referenceType + referenceId
InventoryCheck existing default location by merchantId
LedgerCheck existing pending job by ledgerId

8.2. Error Handling

StrategyUsed By
Skip + LogFinance (logs error, continues processing)
DLQ (Dead Letter Queue)Search CDC (nx.bana.dlq.cdc)
Retry via APILedger (manual retry endpoint)
Startup RecoveryLedger (re-enqueues stalled jobs)

8.3. Serialization

DirectionFormat
ProducerJSON serializer (values), string serializer (keys)
ConsumerJSON deserializer (values), string deserializer (keys)

9. Infrastructure

9.1. Development (Docker Compose)

3-node Kafka cluster:

  • nx-kafka-1:29092
  • nx-kafka-2:29092
  • nx-kafka-3:29092

9.2. Environment Variables

VariableDescription
APP_ENV_KAFKA_BROKERSComma-separated broker addresses
APP_ENV_KAFKA_CLIENT_IDService-specific client identifier
APP_ENV_KAFKA_GROUP_IDConsumer group identifier
APP_ENV_KAFKA_CONSUMER_COUNTNumber of parallel consumers (Ledger)
APP_ENV_KAFKA_SASL_MECHANISMSASL mechanism (e.g., SCRAM-SHA-512)
APP_ENV_KAFKA_SASL_USERNAMESASL username
APP_ENV_KAFKA_SASL_PASSWORDSASL password
DocumentDescription
CDC / DebeziumChange Data Capture via Kafka
Finance ServiceFinance event consumers
Inventory ServiceInventory event consumers
Ledger ServiceLedger generation pipeline
Search ServiceCDC consumer for Typesense sync

Proprietary and Confidential. Unauthorized copying, distribution, or use of this software is strictly prohibited.