feat: Complete VS-03 DOMAIN - Market Data Ingestion (Batch 2 - 3/7)

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>
This commit is contained in:
2026-08-05 21:16:29 +09:00
parent 85e63cbc83
commit f680579134
4 changed files with 883 additions and 0 deletions
@@ -0,0 +1,136 @@
# VS-03: Market Data Ingestion - Vertical Slice Specification
**Slice ID:** VS-03
**Batch:** 2 (depends on VS-00, VS-02, which are complete)
**Status:** 📋 SPECIFICATION
**Created:** 2026-08-05
---
## Executive Summary
Establish **Market Data Ingestion** system that pulls stock prices, indices, and financial data from external sources (KRX, OpenDart) and normalizes them for downstream signal generation.
**User Goal:** Automated, daily market data collection from Korean exchanges with minimal latency and maximum reliability.
**Non-Goal:**
- Real-time tick data (use Bloomberg/Refinitiv for that)
- Cryptocurrency data
- Forex integration
---
## Acceptance Criteria
### 1. Data Sources ✅
- **KRX OpenAPI:** Stock prices, indices, trading volumes
- **OpenDart API:** Financial statements, disclosure documents
- **Fallback:** Stub data (for testing/demo)
### 2. Data Model ✅
- **Market Daily (PIT):** Date, symbol, open, high, low, close, volume
- **Indices:** KRX 200, KOSPI, KOSDAQ snapshots
- **Company Info:** Sector, industry classification, listing status
### 3. Ingestion Pipeline ✅
- **Schedule:** Daily 9:00 KST (before market open)
- **Retry:** Exponential backoff (3 attempts)
- **Validation:** Schema conformance, duplicate detection
- **Idempotency:** By date + symbol (upsert)
- **Audit:** Correlation ID, row count, error logs
### 4. API Contracts ✅
**Endpoint: POST /api/market/ingest**
```
Request: { dataSource: "KRX|OpenDart", fromDate: "2026-01-01", toDate: "2026-12-31" }
Response: 202 Accepted { jobId, expectedRowCount, status }
```
**Endpoint: GET /api/market/ingest/{jobId}**
```
Response: 200 { status, rowsProcessed, rowsFailed, completedAt }
```
### 5. Data Quality Checks ✅
- No NULL prices (OHLCV)
- Volume >= 0
- High >= Low >= Open >= Close (within reason)
- No future dates
- Deduplication by (date, symbol)
---
## Failure Modes & Recovery
| Scenario | Expected | Recovery |
|----------|----------|----------|
| API timeout | 503, retry in 30s | Auto-retry, exponential backoff |
| Bad data format | DQ quarantine | Manual review, adjust parser |
| Duplicate rows | Idempotent upsert | No effect (already stored) |
| Partial ingestion | Rollback, log error | Retry entire day's batch |
---
## Performance SLAs
| Metric | Target |
|--------|--------|
| Daily ingestion latency | <60 seconds |
| Data freshness | <= 1 trading day old |
| Availability | 99.5% (allow 1 failure/week) |
| Max rows/day | 100,000 (stocks + indices) |
---
## Dependencies
### Inbound (Blocked By)
-**VS-00:** Platform foundation (complete)
-**VS-02:** Permission model (complete)
### Outbound (Unblocks)
- 🔄 **VS-04:** Trade Execution (uses VS-03's price data)
- 🔄 **VS-05:** Signal Generation (consumes VS-03 data)
- 🔄 **VS-06:** Portfolio Optimization (requires clean price history)
---
## Component Breakdown (7 items)
| Component | Status |
|-----------|--------|
| **GOV** | 📋 This spec |
| **DATA** | ⏳ Next: PIT schema |
| **DOMAIN** | ⏳ Data validation + normalization |
| **BE** | ⏳ Ingestion API |
| **ASYNC** | ⏳ Hangfire scheduler + event publishing |
| **FE** | ⏳ Ingestion status dashboard |
| **TESTOPS** | ⏳ Data quality tests |
**Total Duration:** ~6 hours (wall-clock 1 day)
---
## Branching Strategy
All work on `Phase-2-Batch-2` branch, squash to main.
**Commits:**
1. GOV + DATA (spec + contract)
2. DOMAIN (validation logic)
3. BE + ASYNC (API + scheduler)
4. FE + TESTOPS (dashboard + tests)
---
## Sign-Off
| Role | Status | Date |
|------|--------|------|
| Architect | ✅ Draft | 2026-08-05 |
| Data Quality | ⏳ Review | TBD |
+260
View File
@@ -0,0 +1,260 @@
# 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