Database Design
Event Store · Read Models · Time-series
Thiết kế schema đầy đủ cho toàn bộ hệ thống ERP/MES — từ Event Store (Write side), Read Model Projections, đến sensor time-series và Redis cache layer.
Chiến lược Database — Write vs Read
Theo kiến trúc CQRS, hệ thống tách biệt hoàn toàn Write side (Event Store) và Read side (Projections). Mỗi loại dữ liệu dùng engine phù hợp nhất với đặc tính của nó.
Với quy mô Phase 1–2, một PostgreSQL cluster (primary + read replica) xử lý được cả Event Store, Read Models và Sensor OEE (range-partitioned). Chỉ cần tách ra database riêng khi đo được bottleneck thực tế — tránh phụ thuộc vào extension bên ngoài (không dùng TimescaleDB).
Event Store Schema — PostgreSQL
Event Store là trung tâm của Write side. Ba bảng chính: events (nguồn sự thật), snapshots (tối ưu replay), outbox (đảm bảo publish to Kafka không mất).
DDL — Bảng events
CREATE TABLE events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
sequence_number BIGINT NOT NULL,
event_type VARCHAR(200) NOT NULL,
payload JSONB NOT NULL,
metadata JSONB NOT NULL DEFAULT '{}',
correlation_id VARCHAR(100),
causation_id VARCHAR(100),
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
occurred_at TIMESTAMPTZ NOT NULL,
recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
-- Optimistic locking: no two events same aggregate + sequence
CONSTRAINT uq_events_aggregate_seq
UNIQUE (aggregate_type, aggregate_id, sequence_number)
);
-- Indexes
CREATE INDEX idx_events_aggregate
ON events (aggregate_type, aggregate_id, sequence_number);
CREATE INDEX idx_events_type
ON events (event_type, occurred_at DESC);
CREATE INDEX idx_events_correlation
ON events (correlation_id)
WHERE correlation_id IS NOT NULL;
CREATE INDEX idx_events_tenant
ON events (tenant_id, recorded_at DESC);
-- Partition by recorded_at for large-scale (optional Phase 3)
-- PARTITION BY RANGE (recorded_at)
DDL — Bảng snapshots
CREATE TABLE snapshots (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
at_sequence BIGINT NOT NULL,
state JSONB NOT NULL,
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
-- Chỉ giữ snapshot mới nhất mỗi aggregate
CONSTRAINT uq_snapshots_aggregate
UNIQUE (aggregate_type, aggregate_id)
);
CREATE INDEX idx_snapshots_aggregate
ON snapshots (aggregate_type, aggregate_id);
-- Trigger snapshot sau mỗi 500 events (thực hiện ở Application layer)
DDL — Bảng outbox (Transactional Outbox Pattern)
CREATE TABLE outbox (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
event_type VARCHAR(200) NOT NULL,
payload JSONB NOT NULL,
kafka_topic VARCHAR(200) NOT NULL,
kafka_key VARCHAR(200) NOT NULL,
-- PENDING | PUBLISHED | FAILED
status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
retry_count INT NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
published_at TIMESTAMPTZ
);
CREATE INDEX idx_outbox_pending
ON outbox (created_at)
WHERE status = 'PENDING';
-- Polling interval: 100ms. Debezium CDC là alternative tốt hơn.
Khi 2 requests cùng lúc ghi vào aggregate A, UNIQUE (aggregate_type, aggregate_id, sequence_number) đảm bảo chỉ 1 thành công. Request thứ 2 nhận UniqueViolationException → retry từ đầu (load aggregate + replay + reapply command).
Read Model — Production Module
Projection tables được tạo và cập nhật bởi Kafka consumers khi nhận Domain Events từ Event Store. Đây là Read-only — không bao giờ ghi trực tiếp vào đây từ Command side.
DDL — production_orders (Projection)
CREATE TABLE production_orders (
id UUID PRIMARY KEY,
order_number VARCHAR(50) NOT NULL UNIQUE,
product_id UUID NOT NULL,
product_code VARCHAR(100) NOT NULL,
planned_qty INT NOT NULL,
actual_qty INT NOT NULL DEFAULT 0,
scrap_qty INT NOT NULL DEFAULT 0,
-- PLANNED | IN_PROGRESS | ON_HOLD | COMPLETED | CANCELLED
status VARCHAR(30) NOT NULL DEFAULT 'PLANNED',
machine_id UUID,
machine_code VARCHAR(50),
operator_id UUID,
operator_name VARCHAR(100),
shift VARCHAR(10), -- MORNING | AFTERNOON | NIGHT
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
planned_start TIMESTAMPTZ,
actual_start TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_prod_orders_status ON production_orders (status, planned_start DESC);
CREATE INDEX idx_prod_orders_machine ON production_orders (machine_id, work_date);
CREATE INDEX idx_prod_orders_tenant ON production_orders (tenant_id, updated_at DESC);
Read Model — Quality & Maintenance Modules
Bảng schemas quan trọng
| Bảng | Mục đích | Trigger event | Index chính |
|---|---|---|---|
| quality_inspections | Header phiếu kiểm tra | QualityInspectionCreated | status, order_id, inspected_at |
| inspection_results | Chi tiết từng thông số đo | InspectionResultRecorded | inspection_id, is_passed |
| non_conformances | Phiếu NCR — sản phẩm không đạt | NonConformanceRaised | status, severity, raised_at |
| assets | Danh mục thiết bị / máy móc | AssetRegistered | asset_type, status, next_maintenance |
| maintenance_work_orders | Lệnh bảo trì (PM/CM/PdM) | WorkOrderCreated | asset_id, status, scheduled_at |
Read Model — Inventory & Finance Modules
DDL — stock_levels (điểm quan trọng)
CREATE TABLE stock_levels (
id UUID PRIMARY KEY,
material_id UUID NOT NULL,
material_code VARCHAR(100) NOT NULL,
material_name VARCHAR(200) NOT NULL,
location_id UUID NOT NULL REFERENCES inventory_locations(id),
on_hand_qty DECIMAL(15,4) NOT NULL DEFAULT 0,
reserved_qty DECIMAL(15,4) NOT NULL DEFAULT 0,
-- available = on_hand - reserved
available_qty DECIMAL(15,4) GENERATED ALWAYS AS
(on_hand_qty - reserved_qty) STORED,
reorder_point DECIMAL(15,4) NOT NULL DEFAULT 0,
unit VARCHAR(20) NOT NULL,
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT uq_stock_material_location
UNIQUE (material_id, location_id, tenant_id)
);
-- Cảnh báo reorder: query materialized view hàng giờ
CREATE MATERIALIZED VIEW mv_reorder_alerts AS
SELECT * FROM stock_levels
WHERE available_qty <= reorder_point
AND reorder_point > 0;
CREATE UNIQUE INDEX ON mv_reorder_alerts (id);
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_reorder_alerts;
Sensor OEE — PostgreSQL Range-Partitioned Tables
Dữ liệu telemetry từ PLC/edge được ingest vào PostgreSQL plain tables với RANGE partitioning theo tháng — không cần TimescaleDB extension. OEE A×P×Q được tính thuần C# (OeeFormulas.cs) và upsert vào snapshot mỗi khi trigger.
DDL — sensor_readings (migration 022)
CREATE TABLE sensor_readings (
id BIGSERIAL,
machine_code VARCHAR(50) NOT NULL,
reading_type VARCHAR(50) NOT NULL, -- cycle_count | good_count | reject_count | cycle_time_ms | power_kw | fault_code
value NUMERIC(18,4) NOT NULL,
unit VARCHAR(20),
source VARCHAR(100),
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
) PARTITION BY RANGE (recorded_at);
CREATE TABLE machine_uptime_events (
id BIGSERIAL,
machine_code VARCHAR(50) NOT NULL,
event_type VARCHAR(30) NOT NULL, -- start | stop | fault | resume | planned_stop
reason TEXT,
source VARCHAR(100),
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
) PARTITION BY RANGE (recorded_at);
-- Monthly partitions created by DO $$ block in migration (13 months)
-- Example:
CREATE TABLE sensor_readings_2026_07
PARTITION OF sensor_readings
FOR VALUES FROM ('2026-07-01') TO ('2026-08-01');
DDL — oee_daily_snapshots
CREATE TABLE oee_daily_snapshots (
id BIGSERIAL PRIMARY KEY,
machine_code VARCHAR(50) NOT NULL,
snapshot_date DATE NOT NULL,
shift_code VARCHAR(20) NOT NULL DEFAULT 'Morning',
planned_minutes NUMERIC(10,2) NOT NULL DEFAULT 480,
uptime_minutes NUMERIC(10,2) NOT NULL DEFAULT 0,
downtime_minutes NUMERIC(10,2) NOT NULL DEFAULT 0,
cycle_count BIGINT NOT NULL DEFAULT 0,
good_count BIGINT NOT NULL DEFAULT 0,
reject_count BIGINT NOT NULL DEFAULT 0,
ideal_cycle_sec NUMERIC(10,4),
-- Generated columns (auto-computed by DB)
oee_availability NUMERIC(6,2) GENERATED ALWAYS AS
(CASE WHEN planned_minutes > 0 THEN ROUND(uptime_minutes / planned_minutes * 100, 2) END) STORED,
oee_quality NUMERIC(6,2) GENERATED ALWAYS AS
(CASE WHEN cycle_count > 0 THEN ROUND(good_count::NUMERIC / cycle_count * 100, 2) END) STORED,
-- Stored columns (computed by OeeFormulas.cs + upserted)
oee_performance NUMERIC(6,2),
oee_total NUMERIC(6,2),
tenant_id VARCHAR(50) NOT NULL DEFAULT 'default',
calculated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT uq_oee_snapshot UNIQUE (machine_code, snapshot_date, shift_code, tenant_id)
);
Sổ cái và công nợ — migrations 031–035
Phần kế toán là nơi giá vốn (Inventory) và giá thành (Finance) chảy vào sổ sách. Ba quyết định schema ở đây đều xuất phát từ một nguyên tắc: số liệu tài chính sai trong im lặng nguy hiểm hơn nhiều so với một lỗi nổ ra ồn ào.
DDL — hệ thống tài khoản và kỳ kế toán
CREATE TABLE chart_of_accounts (
id UUID PRIMARY KEY,
account_number VARCHAR(12) NOT NULL,
name VARCHAR(200) NOT NULL,
account_type VARCHAR(20) NOT NULL, -- Asset|Liability|Equity|Revenue|Expense
parent_number VARCHAR(12),
is_postable BOOLEAN NOT NULL DEFAULT TRUE,
is_active BOOLEAN NOT NULL DEFAULT TRUE,
tenant_id VARCHAR(50) NOT NULL,
-- Cây tài khoản dựa trên TIỀN TỐ số hiệu: 111 -> 1111. Gán cha tuỳ ý làm việc
-- cộng dồn lên cấp trên ra số vô nghĩa mà bảng cân đối vẫn cân — không ai phát hiện.
CONSTRAINT chart_of_accounts_parent_prefix_check
CHECK (parent_number IS NULL
OR (length(account_number) > length(parent_number)
AND account_number LIKE parent_number || '%'))
);
CREATE TABLE accounting_periods (
id UUID PRIMARY KEY,
year INTEGER NOT NULL,
month INTEGER NOT NULL,
fiscal_year_start DATE NOT NULL, -- số dư đầu kỳ của TK doanh thu/chi phí tính từ đây
is_closed BOOLEAN NOT NULL DEFAULT FALSE,
closed_at TIMESTAMPTZ,
closed_by VARCHAR(100),
tenant_id VARCHAR(50) NOT NULL
);
DDL — gl_postings, bảng ghi sổ chỉ-thêm (điểm quan trọng)
Tách gl_postings khỏi journal_entry_lines là quyết định
có chủ đích. Bút toán bị đảo vẫn nằm trên sổ — chỉ được đánh dấu là đã đảo, và có
thêm một bút toán ngược lại để triệt tiêu. Nếu báo cáo đọc dòng chứng từ rồi lọc
status = 'Posted' thì bút toán gốc biến mất trong khi bút toán đảo vẫn còn:
tự tay làm lệch sổ. Đọc từ bảng này thì không có bộ lọc nào để mà đặt sai.
CREATE TABLE gl_postings (
entry_id UUID NOT NULL,
line_number INTEGER NOT NULL,
entry_number VARCHAR(30) NOT NULL,
posting_date DATE NOT NULL,
account_number VARCHAR(12) NOT NULL,
debit NUMERIC(20,2) NOT NULL DEFAULT 0,
credit NUMERIC(20,2) NOT NULL DEFAULT 0,
cost_center VARCHAR(50),
source VARCHAR(20) NOT NULL,
tenant_id VARCHAR(50) NOT NULL,
PRIMARY KEY (entry_id, line_number),
-- Một dòng chỉ ghi MỘT bên; cho cả hai bên thì mọi phép cộng theo bên mất ý nghĩa.
CONSTRAINT gl_posting_one_side_check
CHECK (debit >= 0 AND credit >= 0 AND NOT (debit > 0 AND credit > 0))
);
DDL — chống ghi trùng bút toán tự động
Kafka giao ít nhất một lần: cùng một sự kiện có thể được xử lý lại sau lỗi hoặc sau khi consumer khởi động lại. Ghi hai lần cùng một bút toán là nhân đôi số liệu kế toán. Khoá đặt ở DB chứ không kiểm trong code — kiểm-rồi-ghi vẫn hở khi hai tiến trình chạy song song.
CREATE UNIQUE INDEX ux_journal_auto_reference
ON journal_entries (tenant_id, reference_type, reference_id)
WHERE source <> 'Manual' AND reference_type IS NOT NULL;
-- Nghiệp vụ KHÔNG vào được sổ để lại một dòng ở đây kèm lý do, có đường ghi lại.
-- Sổ sách thiếu một nghiệp vụ mà không ai biết hỏng nặng hơn nhiều so với
-- một hàng đợi lỗi cần người xem.
CREATE TABLE ledger_posting_failures (
id UUID PRIMARY KEY,
rule_code VARCHAR(40) NOT NULL,
reference_type VARCHAR(50) NOT NULL,
reference_id VARCHAR(100) NOT NULL,
amount NUMERIC(20,2) NOT NULL DEFAULT 0,
error_code VARCHAR(50) NOT NULL,
error_message TEXT NOT NULL,
payload JSONB NOT NULL,
resolved_at TIMESTAMPTZ,
tenant_id VARCHAR(50) NOT NULL
);
Bảng công nợ hai chiều — migrations 034 & 035
| Bảng | Vai trò | Ràng buộc đáng chú ý |
|---|---|---|
ar_invoices |
Hoá đơn bán — doanh thu ghi nhận tại đây, không phải lúc giao hàng | outstanding_amount là cột sinh gross − paid, không thể lệch do quên cập nhật |
ar_receipts + ar_receipt_allocations |
Phiếu thu và phần gán vào từng hoá đơn | allocated_amount <= amount — thu dư là tiền trả trước hợp lệ, gán dư thì không |
ap_invoices |
Hoá đơn nhà cung cấp, lưu riêng số hoá đơn của NCC | Unique (tenant, supplier_code, supplier_invoice_number) — nhận trùng là trả tiền hai lần cho một lần mua |
ap_payments + ap_payment_allocations |
Phiếu chi và phần gán | Phải gán HẾT: tiền ra mà không gắn khoản nợ nào là tiền không rõ đi đâu |
Tài khoản trung gian 3388 — hàng về chưa có hoá đơn
Nhập kho ghi Nợ 152 / Có 3388 chứ không ghi thẳng vào 331: lúc nhập kho
thì chưa có hoá đơn nhà cung cấp — chưa biết số hoá đơn, chưa có thuế đầu vào, số
tiền có thể còn lệch. Hoá đơn về mới Nợ 3388 + Nợ 133 / Có 331. Nhờ vậy số dư
3388 luôn bằng đúng giá trị hàng đã nhận mà chưa nhận hoá đơn, và giải trình được theo
từng đơn mua.
Thứ tự sự kiện trong outbox — migration 032
Lỗi đã gặp thật: OutboxPublisher sắp xếp theo created_at, mà
NOW() của Postgres là thời điểm bắt đầu transaction — mọi sự kiện ghi
trong cùng một transaction có created_at giống hệt nhau, thứ tự publish thành
ngẫu nhiên. Trước phần kế toán không BC nào lộ ra vì mỗi request chỉ sinh một sự kiện.
Bút toán dựng-và-ghi-sổ sinh 6 sự kiện một lượt và projection nhận
Posted trước khi đủ dòng.
ALTER TABLE outbox ADD COLUMN seq BIGSERIAL;
DROP INDEX idx_outbox_pending;
CREATE INDEX idx_outbox_pending ON outbox (seq) WHERE status = 'PENDING';
-- Publisher đổi ORDER BY created_at -> ORDER BY seq
Redis 7 — Cache & Hot Data
Redis cache những gì được đọc nhiều nhất và không cần tính nhất quán tức thì. Mọi cache đều có TTL và chiến lược invalidation rõ ràng khi Projection cập nhật.
Key Naming Convention
| Pattern | Data type | TTL | Invalidation |
|---|---|---|---|
session:{userId} | Hash | 8 giờ | Logout / password change |
token:blacklist:{jti} | String | Token expiry | Tự hết hạn |
dashboard:production:{tenant}:{date} | Hash | 60 giây | Projection update event |
machine:status:{machineId} | Hash | 10 giây | MachineStatusChanged event |
product:{productId} | Hash | 1 giờ | ProductUpdated event |
machine:{machineId} | Hash | 1 giờ | MachineUpdated event |
ratelimit:{userId}:{endpoint} | Counter (INCR) | 60 giây | Tự hết hạn |
lock:order:{orderId} | String (SET NX) | 30 giây | Release sau xử lý |
Khi Kafka consumer cập nhật Projection (PostgreSQL Read Model), nó đồng thời xóa/cập nhật Redis key tương ứng. Pattern: DEL dashboard:production:{tenant}:* sau khi Production Order thay đổi trạng thái. Tránh dùng time-based invalidation thuần túy cho dữ liệu quan trọng.
Tài liệu tiếp theo
Domain Events Catalog — Commands, Events & Payloads
Xem tài liệu → ← Về trang chủ