Skip to content

Kafka Event Streaming

1. Tổng quan

BANA dùng Apache Kafka cho giao tiếp bất đồng bộ, hướng sự kiện giữa các microservice. Sự kiện được publish khi có thay đổi trạng thái nghiệp vụ (ví dụ thanh toán thành công, tạo merchant) và được consume bởi các dịch vụ phía dưới phản ứng với những thay đổi đó.

Tất cả tên topic được định nghĩa tập trung trong @nx/core tại src/queues/kafka/topics.ts.

2. Danh mục Topic

Quy ước tên topic: nx.bana.<evt|cmd|cdc|dlq>.<domain>[.<entity>].<event>. Registry sự kiện ứng dụng (ServiceTopicDefs trong packages/core/src/common/kafka/registry.ts):

TopicLoạiMô tả
nx.bana.evt.payment.succeededevtThanh toán bán hàng hoàn tất
nx.bana.evt.inventory.purchase-order.receivedevtĐơn mua được nhận vào kho
nx.bana.evt.inventory.stock.issued-for-saleevtXuất kho cho bán hàng
nx.bana.evt.inventory.stock.adjustedevtĐiều chỉnh tồn kho
nx.bana.evt.kitchen.ticket-item.status-changedevtTrạng thái mục phiếu bếp thay đổi
nx.bana.evt.inventory.material.transferredevtChuyển nguyên vật liệu giữa các kho
nx.bana.evt.inventory.material.stock-changedevtTồn nguyên vật liệu thay đổi
nx.bana.evt.signal.activity.notifiedevtThông báo activity (→ WebSocket)
nx.bana.cmd.ledger.generatecmdTạo tài liệu ledger (nội bộ ledger)

Thay đổi merchant, product và product-variant không phải evt topic - chúng lan truyền dưới dạng CDC Debezium nx.bana.cdc.<schema>.<Table> (xem CDC). Không có topic commerce.initialized.

3. Ma trận Producer & Consumer

TopicProducerConsumer
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 (chuyển kho)
nx.bana.evt.inventory.material.stock-changedinventory(chưa có consumer)
nx.bana.evt.signal.activity.notified(chưa có producer)signal
nx.bana.cmd.ledger.generateledgerledger (nội bộ)

4. Sơ đồ Luồng Message

5. Consumer Group

Dịch vụClient 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. Cấu hình Consumer

6.1. Consumer Finance

Cài đặtGiá trị
Topicnx.bana.evt.payment.succeeded, …inventory.purchase-order.received, …inventory.stock.issued-for-sale, …inventory.stock.adjusted, CDC public.Merchant
Auto-commitbật
Fallback Modeearliest
DeserializerKey dạng String, Value dạng JSON

6.2. Consumer Inventory

Cài đặtGiá trị
Topicnx.bana.evt.payment.succeeded, …kitchen.ticket-item.status-changed, …inventory.material.transferred, CDC public.Merchant, CDC public.ProductVariant
Fallback Modeearliest

6.3. Consumer Ledger

Cài đặtGiá trị
Topicnx.bana.cmd.ledger.generate
Auto-commitfalse (commit thủ công sau khi xử lý)
Số lượng ConsumerCó thể cấu hình qua APP_ENV_KAFKA_CONSUMER_COUNT
Acks (Producer)ALL
Idempotent (Producer)true

6.4. Consumer Search CDC

Cài đặtGiá trị
Topic27 topic CDC Debezium
Auto-commitfalse (quản lý offset thủ công)
Max Batch Size200
Flush Interval2000ms
Max Wait Time500ms
Max Bytes5MB

Xem CDC / Debezium để biết chi tiết.

7. Ví dụ Payload Sự kiện

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. Thực hành Tốt nhất

8.1. Idempotency

Tất cả consumer đều thực hiện kiểm tra idempotency để xử lý an toàn các message trùng lặp:

Dịch vụChiến lược
FinanceKiểm tra giao dịch hiện có theo merchantId + referenceType + referenceId
InventoryKiểm tra vị trí mặc định hiện có theo merchantId
LedgerKiểm tra tác vụ đang chờ hiện có theo ledgerId

8.2. Xử lý Lỗi

Chiến lượcĐược dùng bởi
Bỏ qua + LogFinance (log lỗi, tiếp tục xử lý)
DLQ (Dead Letter Queue)Search CDC (nx.bana.dlq.cdc)
Retry qua APILedger (endpoint retry thủ công)
Phục hồi khi khởi độngLedger (đưa lại các tác vụ bị treo vào hàng đợi)

8.3. Serialization

HướngĐịnh dạng
ProducerJSON serializer (value), string serializer (key)
ConsumerJSON deserializer (value), string deserializer (key)

9. Hạ tầng

9.1. Phát triển (Docker Compose)

Cluster Kafka 3 node:

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

9.2. Biến Môi trường

BiếnMô tả
APP_ENV_KAFKA_BROKERSĐịa chỉ broker cách nhau bằng dấu phẩy
APP_ENV_KAFKA_CLIENT_IDĐịnh danh client riêng cho dịch vụ
APP_ENV_KAFKA_GROUP_IDĐịnh danh consumer group
APP_ENV_KAFKA_CONSUMER_COUNTSố consumer song song (Ledger)
APP_ENV_KAFKA_SASL_MECHANISMCơ chế SASL (ví dụ SCRAM-SHA-256)
APP_ENV_KAFKA_SASL_USERNAMEUsername SASL
APP_ENV_KAFKA_SASL_PASSWORDPassword SASL

10. Tài liệu Liên quan

Tài liệuMô tả
CDC / DebeziumChange Data Capture qua Kafka
Finance ServiceConsumer sự kiện Finance
Inventory ServiceConsumer sự kiện Inventory
Ledger ServicePipeline tạo Ledger
Search ServiceConsumer CDC để đồng bộ Typesense

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