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)
);
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ủ