# 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