f680579134
Implements market data validation and normalization: ✅ GOV: Market data ingestion specification - KRX/OpenDart data sources - Daily scheduling (9:00 KST) - Quality SLAs (99.5% availability) ✅ DATA: PIT-compliant schema (4 tables) - daily_prices: OHLCV with versioning - indices: Market indices snapshots - companies: Master data - ingestion_jobs: Audit trail ✅ DOMAIN: Policy logic (12 tests, 12/12 PASS) - ValidatePrice: OHLC constraints, date checks - IsDuplicate: Prevent redundant entries - NormalizePrice: Rounding, filtering - ClassifyQualityIssue: Quality scoring (0-100) - ValidateBatch: Aggregate metrics AGENTS.md v16.0 compliance: ✅ Necessity: WBS Phase 2 Batch 2 ✅ Simplicity: Pure validation logic, no I/O ✅ Idempotency: By (symbol, trading_date) ✅ Safety: Immutable history with versioning ✅ Quality gates: Data quality scoring Phase 2 Progress: 1/4 Batches (VS-03 GOV+DATA+DOMAIN COMPLETE) Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
261 lines
7.2 KiB
Markdown
261 lines
7.2 KiB
Markdown
# VS-03: Market Data Ingestion - Data Contract
|
|
|
|
**Slice ID:** VS-03
|
|
**Phase:** Data Layer (write model)
|
|
**Status:** Specification Ready
|
|
|
|
---
|
|
|
|
## Write Model (Normalized, 3NF)
|
|
|
|
### Table: `market_data.daily_prices` (Core)
|
|
|
|
```sql
|
|
CREATE TABLE market_data.daily_prices (
|
|
-- Identity
|
|
price_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
symbol VARCHAR(20) NOT NULL,
|
|
trading_date DATE NOT NULL,
|
|
|
|
-- OHLCV
|
|
open_price DECIMAL(10, 2) NOT NULL CHECK (open_price > 0),
|
|
high_price DECIMAL(10, 2) NOT NULL CHECK (high_price > 0),
|
|
low_price DECIMAL(10, 2) NOT NULL CHECK (low_price > 0),
|
|
close_price DECIMAL(10, 2) NOT NULL CHECK (close_price > 0),
|
|
adjusted_close DECIMAL(10, 2),
|
|
volume BIGINT NOT NULL CHECK (volume >= 0),
|
|
|
|
-- PIT (Point-in-Time) Compliance
|
|
published_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
revision INT NOT NULL DEFAULT 1,
|
|
|
|
-- Audit
|
|
data_source VARCHAR(50) NOT NULL, -- 'KRX', 'OpenDart', 'Stub'
|
|
ingestion_job_id UUID,
|
|
correlation_id UUID,
|
|
|
|
-- Soft-delete (never delete, only version)
|
|
removed_at TIMESTAMP,
|
|
|
|
CONSTRAINT unique_daily_price UNIQUE (symbol, trading_date, revision),
|
|
CONSTRAINT valid_prices CHECK (low_price <= open_price AND open_price <= high_price)
|
|
);
|
|
|
|
CREATE INDEX idx_daily_prices_symbol_date ON market_data.daily_prices(symbol, trading_date DESC);
|
|
CREATE INDEX idx_daily_prices_published ON market_data.daily_prices(published_at DESC);
|
|
```
|
|
|
|
### Table: `market_data.indices` (Supplementary)
|
|
|
|
```sql
|
|
CREATE TABLE market_data.indices (
|
|
-- Identity
|
|
index_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
index_code VARCHAR(20) NOT NULL, -- 'KOSPI', 'KRX200', 'KOSDAQ'
|
|
trading_date DATE NOT NULL,
|
|
|
|
-- OHLCV
|
|
open_value DECIMAL(10, 2) NOT NULL,
|
|
high_value DECIMAL(10, 2) NOT NULL,
|
|
low_value DECIMAL(10, 2) NOT NULL,
|
|
close_value DECIMAL(10, 2) NOT NULL,
|
|
change_percent DECIMAL(5, 2),
|
|
volume BIGINT,
|
|
|
|
-- PIT
|
|
published_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
revision INT NOT NULL DEFAULT 1,
|
|
|
|
-- Audit
|
|
data_source VARCHAR(50) NOT NULL,
|
|
correlation_id UUID,
|
|
|
|
removed_at TIMESTAMP,
|
|
|
|
CONSTRAINT unique_index UNIQUE (index_code, trading_date, revision)
|
|
);
|
|
|
|
CREATE INDEX idx_indices_code_date ON market_data.indices(index_code, trading_date DESC);
|
|
```
|
|
|
|
### Table: `market_data.companies` (Master)
|
|
|
|
```sql
|
|
CREATE TABLE market_data.companies (
|
|
-- Identity
|
|
company_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
symbol VARCHAR(20) NOT NULL UNIQUE,
|
|
|
|
-- Master Data
|
|
korean_name VARCHAR(100) NOT NULL,
|
|
english_name VARCHAR(100),
|
|
sector VARCHAR(50),
|
|
industry VARCHAR(100),
|
|
listing_date DATE,
|
|
|
|
-- Status
|
|
listing_status VARCHAR(20) NOT NULL DEFAULT 'Active', -- Active, Suspended, Delisted
|
|
market VARCHAR(20) NOT NULL, -- 'KOSPI', 'KOSDAQ', 'KONEX'
|
|
|
|
-- PIT
|
|
published_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
revision INT NOT NULL DEFAULT 1,
|
|
removed_at TIMESTAMP,
|
|
|
|
-- Audit
|
|
last_updated TIMESTAMP,
|
|
data_source VARCHAR(50),
|
|
|
|
CONSTRAINT unique_company UNIQUE (symbol, revision)
|
|
);
|
|
|
|
CREATE INDEX idx_companies_symbol ON market_data.companies(symbol);
|
|
```
|
|
|
|
### Table: `market_data.ingestion_jobs` (Audit)
|
|
|
|
```sql
|
|
CREATE TABLE market_data.ingestion_jobs (
|
|
-- Identity
|
|
job_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
job_run_id UUID NOT NULL, -- Hangfire RunId
|
|
|
|
-- Input
|
|
data_source VARCHAR(50) NOT NULL,
|
|
from_date DATE NOT NULL,
|
|
to_date DATE NOT NULL,
|
|
|
|
-- Progress
|
|
status VARCHAR(50) NOT NULL DEFAULT 'Queued', -- Queued, Running, Completed, Failed
|
|
rows_processed INT DEFAULT 0,
|
|
rows_failed INT DEFAULT 0,
|
|
rows_skipped INT DEFAULT 0,
|
|
|
|
-- Timing
|
|
started_at TIMESTAMP,
|
|
completed_at TIMESTAMP,
|
|
duration_seconds INT,
|
|
|
|
-- Error Handling
|
|
last_error_message TEXT,
|
|
retry_count INT DEFAULT 0,
|
|
|
|
-- Traceability
|
|
correlation_id UUID NOT NULL,
|
|
triggered_by VARCHAR(100), -- 'Scheduler', 'Manual', 'API'
|
|
|
|
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
|
|
CONSTRAINT unique_job_run UNIQUE (job_run_id)
|
|
);
|
|
|
|
CREATE INDEX idx_ingestion_jobs_status ON market_data.ingestion_jobs(status);
|
|
CREATE INDEX idx_ingestion_jobs_dates ON market_data.ingestion_jobs(from_date, to_date);
|
|
```
|
|
|
|
---
|
|
|
|
## Read Model (Denormalized Projections)
|
|
|
|
### View: `market_data.latest_prices` (Cache)
|
|
|
|
```sql
|
|
CREATE VIEW market_data.latest_prices AS
|
|
SELECT DISTINCT ON (symbol)
|
|
symbol,
|
|
trading_date,
|
|
close_price,
|
|
volume,
|
|
published_at
|
|
FROM market_data.daily_prices
|
|
WHERE removed_at IS NULL
|
|
AND published_at <= CURRENT_TIMESTAMP
|
|
ORDER BY symbol, trading_date DESC;
|
|
```
|
|
|
|
---
|
|
|
|
## PIT (Point-in-Time) Query Pattern
|
|
|
|
```sql
|
|
-- Fetch prices as of 2026-06-30
|
|
SELECT symbol, open_price, close_price, volume
|
|
FROM market_data.daily_prices
|
|
WHERE trading_date <= '2026-06-30'
|
|
AND published_at <= '2026-06-30'::timestamp
|
|
AND removed_at IS NULL
|
|
ORDER BY symbol, trading_date DESC
|
|
LIMIT 1 PER symbol;
|
|
```
|
|
|
|
---
|
|
|
|
## Migration Strategy
|
|
|
|
1. **0033_market_data_schema.sql**
|
|
- Create market_data schema
|
|
- Define daily_prices, indices, companies, ingestion_jobs tables
|
|
- Add PK, FK, constraints
|
|
|
|
2. **0034_market_data_indexes.sql**
|
|
- Create performance indexes
|
|
- Partition by year (optional, if 10M+ rows/year)
|
|
|
|
3. **0035_market_data_audit.sql**
|
|
- Create audit trigger (log all writes)
|
|
- Set up row-level security (market access control)
|
|
|
|
---
|
|
|
|
## Data Dictionary
|
|
|
|
| Column | Type | Purpose |
|
|
|--------|------|---------|
|
|
| symbol | VARCHAR(20) | Stock ticker (e.g., '005930' for Samsung) |
|
|
| trading_date | DATE | Market trading date (YYYY-MM-DD) |
|
|
| open_price | DECIMAL(10,2) | Opening price |
|
|
| close_price | DECIMAL(10,2) | Closing price |
|
|
| volume | BIGINT | Trading volume (shares) |
|
|
| published_at | TIMESTAMP | PIT anchor (when row became "true") |
|
|
| revision | INT | Version number (immutable history) |
|
|
| removed_at | TIMESTAMP | Soft-delete marker (NULL = active) |
|
|
| correlation_id | UUID | Trace this data ingestion back to job |
|
|
|
|
---
|
|
|
|
## Idempotency & Upsert Strategy
|
|
|
|
**Idempotency Key:** `(symbol, trading_date)`
|
|
|
|
**Upsert SQL:**
|
|
```sql
|
|
INSERT INTO market_data.daily_prices (symbol, trading_date, open_price, high_price, low_price, close_price, volume, published_at, revision, correlation_id, data_source)
|
|
VALUES (@symbol, @date, @open, @high, @low, @close, @volume, CURRENT_TIMESTAMP, 1, @corrId, @source)
|
|
ON CONFLICT (symbol, trading_date, revision) DO UPDATE SET
|
|
open_price = EXCLUDED.open_price,
|
|
close_price = EXCLUDED.close_price,
|
|
volume = EXCLUDED.volume,
|
|
published_at = CURRENT_TIMESTAMP,
|
|
revision = market_data.daily_prices.revision + 1
|
|
WHERE EXCLUDED.published_at > market_data.daily_prices.published_at;
|
|
```
|
|
|
|
**Effect:** Same-day re-ingestion updates the row; older data is immutable (PIT principle).
|
|
|
|
---
|
|
|
|
## Testing & Validation
|
|
|
|
**Unit Tests (SQL):**
|
|
- Constraints enforced (negative prices rejected)
|
|
- Unique keys prevent duplicates
|
|
- Soft-delete preserves history
|
|
- PIT query returns correct version
|
|
|
|
**Integration Tests:**
|
|
- Ingest 100 rows, verify count
|
|
- Duplicate ingestion (same date/symbol) increments revision
|
|
- Upsert with newer timestamp overwrites
|
|
|