API Events
@nx/searchlà một CDC consumer. Nó đăng ký 18 Debezium topic và chỉ phát message tới một dead-letter queue. Không có business event outbound, không WebSocket, không BullMQ. Tên topic theo quy ước Debezium{prefix}.{schema}.{table}với prefixnx.seller.
1. Inbound - Kafka (CDC)
Tất cả topic được định nghĩa trong SearchCollections.CDC (src/common/kafka-topics.ts); consumer đăng ký ALL_CDC_TOPICS. Handler là CDCService.handleBatch() cho mọi topic.
Nguồn document trực tiếp (ánh xạ trong TableToCollectionMap)
| Topic | Schema.Table | → Collection |
|---|---|---|
nx.seller.public.Organizer | public.Organizer | organizers |
nx.seller.public.Merchant | public.Merchant | merchants |
nx.seller.public.Category | public.Category | categories |
nx.seller.public.Device | public.Device | devices |
nx.seller.public.SaleChannel | public.SaleChannel | sale-channels |
nx.seller.public.Product | public.Product | products |
nx.seller.public.ProductInfo | public.ProductInfo | products (i18n partial) |
nx.seller.public.ProductVariant | public.ProductVariant | product-variants |
nx.seller.inventory.InventoryStock | inventory.InventoryStock | inventories |
Nguồn chỉ cascade (không phải nguồn document trực tiếp - lan toả qua CDCCascadeService)
| Topic | Schema.Table | Tính lại |
|---|---|---|
nx.seller.public.ProductCategory | public.ProductCategory | products + product-variants categoryIds |
nx.seller.public.MetaLink | public.MetaLink | product → variant productMetaLinks; variant metaLinks |
nx.seller.pricing.FareSet | pricing.FareSet | variant fareSet + defaultPrice |
nx.seller.pricing.Fare | pricing.Fare | variant fareSet + defaultPrice |
nx.seller.public.ProductBundler | public.ProductBundler | variant comboItems |
nx.seller.inventory.InventoryItem | inventory.InventoryItem | nhóm item của inventories |
nx.seller.inventory.InventoryLocation | inventory.InventoryLocation | nhóm location của inventories |
nx.seller.inventory.InventoryIdentifier | inventory.InventoryIdentifier | inventories identifiers[] |
nx.seller.inventory.Material | inventory.Material | item của inventories (kiểu Material) |
Mã op Debezium
| Op | Ý nghĩa | Hành động |
|---|---|---|
c | tạo | upsert document |
u | cập nhật | upsert document |
d | xoá | xoá document |
r | đọc snapshot | upsert document (nạp lần đầu) |
Bất kỳ op nào mà mapper trả về
null(đã soft-delete) đều trở thành lệnh xoá trên Typesense.
2. Outbound - Kafka
| Topic | Kích hoạt | Bên tiêu thụ | Payload |
|---|---|---|---|
nx.seller.cdc.dlq (mặc định, ghi đè qua APP_ENV_CDC_DLQ_TOPIC) | Một message CDC xử lý thất bại sau khi thử lại | Ops / công cụ replay thủ công | Message CDC gốc + metadata lỗi |
3. Inbound - BullMQ
N/A - không có BullMQ consumer. (Thiết kế embedding-queue trong selection report là một phương án tương lai, chưa được triển khai.)
4. Outbound - BullMQ
N/A - không có BullMQ producer.
5. WebSocket Emissions
N/A - thư viện không phát WebSocket event nào. Cập nhật UI thời gian thực là việc của host service.
6. Payload Schemas
Payload CDC là Debezium envelope, được giải mã bằng avsc. Search không sở hữu các schema này - chúng phản chiếu bảng nguồn. Cấu trúc (rút gọn):
// Debezium change event (Avro-decoded)
interface DebeziumEnvelope<TBefore, TAfter> {
op: 'c' | 'u' | 'd' | 'r';
before: TBefore | null;
after: TAfter | null;
source: {
schema: string; // public | pricing | inventory
table: string; // vd "Product"
lsn: number; // → chặn replay sai thứ tự (SearchVersionFields.SOURCE_LSN)
ts_ms: number;
};
}Tham số truy vấn search (phía consumer, ISearchParams):
interface ISearchParams {
q?: string; // mặc định '*'
limit?: number; // mặc định 10, max 250
offset?: number;
where?: unknown; // kiểu Ignis → filter_by
order?: string | string[]; // "createdAt DESC" → sort_by
include?: IIncludeSpec[]; // bù dữ liệu quan hệ / native join
useCache?: boolean;
cacheTtl?: number;
disableSemanticSearch?: boolean; // bỏ embedding khỏi query_by (thuần từ khoá)
}7. Idempotency & Ordering
| Loại topic | Phân phối | Thứ tự | Phục hồi |
|---|---|---|---|
| Tất cả topic CDC | at-least-once (tắt auto-commit; commit sau batch) | theo key qua Debezium; bộ chặn LSN loại bỏ event cũ | đọc lại offset chưa commit sau khi khởi động lại / circuit breaker đóng |
| DLQ | gửi-rồi-quên | - | replay thủ công |
| Bộ chặn | Trường | Tác dụng |
|---|---|---|
| LSN | source_lsn (SearchVersionFields.SOURCE_LSN) | Bỏ qua event nếu LSN của nó ≤ LSN document đã lưu |
| Tombstone | deleted_at (SearchVersionFields.DELETED_AT) | So sánh để giữ lệnh xoá không bị upsert cũ hoàn tác |