From ca9010b1355137ff1e13708e7da4d293e3da256d Mon Sep 17 00:00:00 2001 From: oosuhada <185910926+oosuhada@users.noreply.github.com> Date: Tue, 8 Sep 2026 19:40:39 +0900 Subject: [PATCH] perf: add measured PostgreSQL query profiling --- .github/workflows/postgres-profile.yml | 69 +++++ README.md | 8 + .../008_operational_query_indexes.sql | 12 + docs/POSTGRESQL_PROFILING.md | 62 ++++ evaluation/postgres_profile.py | 269 ++++++++++++++++++ .../postgres/operational-index-profile.json | 85 ++++++ tests/test_postgres_profile.py | 37 +++ 7 files changed, 542 insertions(+) create mode 100644 .github/workflows/postgres-profile.yml create mode 100644 adapters/postgres/migrations/008_operational_query_indexes.sql create mode 100644 docs/POSTGRESQL_PROFILING.md create mode 100644 evaluation/postgres_profile.py create mode 100644 evidence/postgres/operational-index-profile.json create mode 100644 tests/test_postgres_profile.py diff --git a/.github/workflows/postgres-profile.yml b/.github/workflows/postgres-profile.yml new file mode 100644 index 0000000..4d6ba23 --- /dev/null +++ b/.github/workflows/postgres-profile.yml @@ -0,0 +1,69 @@ +name: postgres-profile + +on: + pull_request: + paths: + - "adapters/postgres/**" + - "evaluation/postgres_profile.py" + - "tests/test_postgres_profile.py" + - ".github/workflows/postgres-profile.yml" + push: + branches: + - main + paths: + - "adapters/postgres/**" + - "evaluation/postgres_profile.py" + - "tests/test_postgres_profile.py" + - ".github/workflows/postgres-profile.yml" + workflow_dispatch: + +permissions: + contents: read + +jobs: + profile: + runs-on: ubuntu-latest + timeout-minutes: 15 + services: + postgres: + image: postgres:16-alpine + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: fabops_bench + ports: + - 5432:5432 + options: >- + --health-cmd "pg_isready -U postgres -d fabops_bench" + --health-interval 5s + --health-timeout 5s + --health-retries 12 + + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Set up Python + uses: actions/setup-python@v5 + with: + python-version: "3.12" + + - name: Install uv + run: python -m pip install --no-cache-dir uv + + - name: Install locked dependencies + run: uv sync --locked --dev + + - name: Lint and unit tests + run: | + uv run ruff check evaluation/postgres_profile.py tests/test_postgres_profile.py + uv run pytest -q tests/test_postgres_profile.py + + - name: Verify PostgreSQL plans and latency direction + run: >- + uv run python -m evaluation.postgres_profile + --dsn postgresql://postgres:postgres@127.0.0.1:5432/fabops_bench + --event-rows 50000 + --case-rows 20000 + --repeats 3 + --strict diff --git a/README.md b/README.md index a2bd925..b23a9cb 100644 --- a/README.md +++ b/README.md @@ -380,6 +380,14 @@ These numbers are **not** PostgreSQL/Neo4j/Redpanda production capacity claims. 이 수치는 PostgreSQL/Neo4j/Redpanda의 **production capacity 주장**이 아닙니다. 자세한 정의와 제한은 `docs/operations/SLO.md`에 있습니다. +### PostgreSQL query-plan evidence / PostgreSQL 실행계획 근거 + +Two real repository query shapes were profiled against an isolated PostgreSQL 16 benchmark database using `EXPLAIN (ANALYZE, BUFFERS)` over 200,000 event rows and 80,000 case rows. A composite `(event_type, sequence DESC)` index reduced the recent-measurement query p50 from **10.230 ms to 0.033 ms** and buffer hits from **12,078 to 14**. A `(classification, lot_id DESC, updated_at DESC)` index reduced the related-case query p50 from **0.162 ms to 0.018 ms**, while removing the previous incremental-sort step. + +실제 repository query shape 두 개를 격리 PostgreSQL 16 환경에서 20만 event / 8만 case fixture와 `EXPLAIN (ANALYZE, BUFFERS)`로 측정했습니다. `(event_type, sequence DESC)` composite index는 recent-measurement query p50을 **10.230 ms → 0.033 ms**, buffer hit을 **12,078 → 14**로 줄였고, `(classification, lot_id DESC, updated_at DESC)` index는 related-case query p50을 **0.162 ms → 0.018 ms**로 줄이면서 기존 incremental-sort 단계를 제거했습니다. + +Full experiment, plan names, variance and limitations: `docs/POSTGRESQL_PROFILING.md` and `evidence/postgres/operational-index-profile.json`. + ### Executed incident exercise / 실제 장애 훈련 The M6 Neo4j dependency outage exercise actually stopped the project Neo4j container, observed degraded readiness, restarted it, and observed recovery: diff --git a/adapters/postgres/migrations/008_operational_query_indexes.sql b/adapters/postgres/migrations/008_operational_query_indexes.sql new file mode 100644 index 0000000..8278a03 --- /dev/null +++ b/adapters/postgres/migrations/008_operational_query_indexes.sql @@ -0,0 +1,12 @@ +-- Indexes selected from measured repository query shapes. +-- Keep these narrow: JSONB payload columns stay in the heap to avoid bloating +-- the write path for the authoritative event/case tables. +BEGIN; + +CREATE INDEX IF NOT EXISTS fabops_event_log_event_type_sequence_idx + ON fabops_event_log(event_type, sequence DESC); + +CREATE INDEX IF NOT EXISTS fabops_cases_classification_lot_updated_idx + ON fabops_cases(classification, lot_id DESC, updated_at DESC); + +COMMIT; diff --git a/docs/POSTGRESQL_PROFILING.md b/docs/POSTGRESQL_PROFILING.md new file mode 100644 index 0000000..2b1878d --- /dev/null +++ b/docs/POSTGRESQL_PROFILING.md @@ -0,0 +1,62 @@ +# PostgreSQL operational query profiling + +## Problem + +Two repository reads are order-sensitive and can become expensive as the event/case tables grow: + +- `recent_measurement_events()` filters one event type and asks for the newest sequences. +- `related_cases()` filters one classification and orders by lot/update recency. + +The existing single-column indexes did not fully match those filter + order shapes. + +## Measurement + +`evaluation/postgres_profile.py` creates an isolated benchmark database, applies the real FabOps migrations, seeds **200,000 events** and **80,000 cases**, and runs the repository-equivalent SQL with `EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON)`. + +Each before/after phase is repeated 15 times. The profiler refuses to seed or drop indexes unless the database name contains `bench` or `test`. + +Command used for the committed result: + +```bash +uv run python -m evaluation.postgres_profile \ + --dsn postgresql://postgres:postgres@127.0.0.1:55432/fabops_bench \ + --event-rows 200000 \ + --case-rows 80000 \ + --repeats 15 +``` + +The DSN above is a disposable local benchmark database, not a production connection string. + +## Change + +Migration `008_operational_query_indexes.sql` adds only the two indexes that match measured repository access patterns: + +```sql +CREATE INDEX fabops_event_log_event_type_sequence_idx + ON fabops_event_log(event_type, sequence DESC); + +CREATE INDEX fabops_cases_classification_lot_updated_idx + ON fabops_cases(classification, lot_id DESC, updated_at DESC); +``` + +No GIN/GiST/partitioning feature was added because these reads do not justify them. + +## Result + +| Query | Before | After | p50 | p95 | Buffer hits | +|---|---|---|---:|---:|---:| +| recent measurements | `fabops_event_log_pkey` | `fabops_event_log_event_type_sequence_idx` | **10.230 → 0.033 ms** | **10.627 → 0.0413 ms** | **12,078 → 14** | +| related cases | `idx_fabops_cases_lot_id` + incremental sort | `fabops_cases_classification_lot_updated_idx` | **0.162 → 0.018 ms** | **0.1799 → 0.0324 ms** | **2,844 → 46** | + +Measured p50 reductions were **99.68%** and **88.89%** respectively on this fixture. + +The complete machine-readable result is committed at `evidence/postgres/operational-index-profile.json`. + +The CI regression gate runs a smaller isolated fixture with `--strict`; it fails if the intended index is not selected, p50 does not improve, or shared buffer hits do not decrease. The smaller CI fixture is a regression direction check, not the source of the README benchmark numbers above. + +## Limitation + +- This is a synthetic portfolio fixture using the real schema/query shape, not a production-fab capacity benchmark. +- Measurements are warm-cache runs on one local PostgreSQL 16 instance. +- The experiment isolates read plans; it does not quantify index write amplification, vacuum behavior, or production concurrency. +- The profiler records latency variance, but it is not a substitute for workload-level contention testing. diff --git a/evaluation/postgres_profile.py b/evaluation/postgres_profile.py new file mode 100644 index 0000000..912e96f --- /dev/null +++ b/evaluation/postgres_profile.py @@ -0,0 +1,269 @@ +from __future__ import annotations + +import argparse +import json +import math +import statistics +from dataclasses import asdict, dataclass +from pathlib import Path +from typing import Any +from urllib.parse import urlparse + +import psycopg + +from adapters.postgres.migrate import apply_migrations + +EVENT_INDEX = "fabops_event_log_event_type_sequence_idx" +CASE_INDEX = "fabops_cases_classification_lot_updated_idx" + + +@dataclass(frozen=True) +class PhaseResult: + p50_ms: float + p95_ms: float + mean_ms: float + variance_ms2: float + plan_nodes: list[str] + index_names: list[str] + shared_hit_blocks: int + shared_read_blocks: int + + +def _percentile(values: list[float], percentile: float) -> float: + ordered = sorted(values) + if not ordered: + return 0.0 + position = (len(ordered) - 1) * percentile + lower = math.floor(position) + upper = math.ceil(position) + if lower == upper: + return ordered[lower] + return ordered[lower] + (ordered[upper] - ordered[lower]) * (position - lower) + + +def _plan_nodes(plan: dict[str, Any]) -> list[str]: + nodes = [str(plan.get("Node Type", "unknown"))] + for child in plan.get("Plans", []): + nodes.extend(_plan_nodes(child)) + return nodes + + +def _index_names(plan: dict[str, Any]) -> list[str]: + names: list[str] = [] + index_name = plan.get("Index Name") + if index_name: + names.append(str(index_name)) + for child in plan.get("Plans", []): + names.extend(_index_names(child)) + return names + + +def _buffer_total(plan: dict[str, Any], key: str) -> int: + total = int(plan.get(key, 0) or 0) + for child in plan.get("Plans", []): + total += _buffer_total(child, key) + return total + + +def _explain_runs(connection: psycopg.Connection[Any], sql: str, params: tuple[Any, ...], repeats: int) -> PhaseResult: + execution_times: list[float] = [] + last_plan: dict[str, Any] = {} + for _ in range(repeats): + row = connection.execute( + f"EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) {sql}", + params, + ).fetchone() + payload = row[0][0] + last_plan = payload["Plan"] + execution_times.append(float(payload["Execution Time"])) + + return PhaseResult( + p50_ms=round(_percentile(execution_times, 0.50), 4), + p95_ms=round(_percentile(execution_times, 0.95), 4), + mean_ms=round(statistics.mean(execution_times), 4), + variance_ms2=round(statistics.pvariance(execution_times), 6), + plan_nodes=_plan_nodes(last_plan), + index_names=_index_names(last_plan), + shared_hit_blocks=_buffer_total(last_plan, "Shared Hit Blocks"), + shared_read_blocks=_buffer_total(last_plan, "Shared Read Blocks"), + ) + + +def _safe_benchmark_database(dsn: str) -> str: + database = urlparse(dsn).path.lstrip("/") + if not database or not any(marker in database.lower() for marker in ("bench", "test")): + raise SystemExit( + "Refusing to seed/drop indexes outside a database whose name contains 'bench' or 'test'." + ) + return database + + +def _seed_fixture(connection: psycopg.Connection[Any], event_rows: int, case_rows: int) -> None: + connection.execute( + "TRUNCATE fabops_decision_audit, fabops_cases, fabops_measurements, fabops_event_log CASCADE" + ) + connection.execute( + """ + INSERT INTO fabops_event_log( + event_id, event_type, event_time, ingested_at, trace_id, lot_id, + equipment_id, chamber_id, schema_version, delivery_status, envelope + ) + SELECT + md5('event-' || g::text)::uuid, + CASE WHEN g <= GREATEST(512, %s / 100) THEN 'process.measurement.recorded.v1' + ELSE 'process.step.completed.v1' END, + now() - (g || ' milliseconds')::interval, + now(), + 'trace-' || (g %% 2000)::text, + 'LOT-' || lpad((g %% 5000)::text, 5, '0'), + 'EQ-' || (g %% 40)::text, + 'CH-' || (g %% 8)::text, + 1, + 'on_time', + jsonb_build_object('event_id', g, 'fixture', true) + FROM generate_series(1, %s) AS g + """, + (event_rows, event_rows), + ) + connection.execute( + """ + INSERT INTO fabops_cases( + case_id, lot_id, classification, state, anomaly_score, + detector_version, case_document, updated_at + ) + SELECT + 'CASE-' || lpad(g::text, 7, '0'), + 'LOT-' || lpad((g %% 5000)::text, 5, '0'), + CASE WHEN g %% 20 = 0 THEN 'physical_excursion' + WHEN g %% 5 = 0 THEN 'sensor_bias_suspected' + ELSE 'data_quality_incident' END, + 'open', + (g %% 100)::double precision / 100.0, + 'profile-fixture-v1', + jsonb_build_object('case_id', g, 'fixture', true), + now() - (g || ' milliseconds')::interval + FROM generate_series(1, %s) AS g + """, + (case_rows,), + ) + connection.execute("ANALYZE fabops_event_log") + connection.execute("ANALYZE fabops_cases") + + +def _profile_one( + connection: psycopg.Connection[Any], + *, + index_name: str, + create_index_sql: str, + query_sql: str, + params: tuple[Any, ...], + repeats: int, +) -> dict[str, Any]: + connection.execute(f'DROP INDEX IF EXISTS "{index_name}"') + connection.execute("DISCARD PLANS") + before = _explain_runs(connection, query_sql, params, repeats) + + connection.execute(create_index_sql) + connection.execute("ANALYZE fabops_event_log") + connection.execute("ANALYZE fabops_cases") + connection.execute("DISCARD PLANS") + after = _explain_runs(connection, query_sql, params, repeats) + + improvement = None + if before.p50_ms > 0: + improvement = round((before.p50_ms - after.p50_ms) / before.p50_ms * 100.0, 2) + return { + "index": index_name, + "before": asdict(before), + "after": asdict(after), + "p50_improvement_percent": improvement, + } + + +def run_profile(dsn: str, *, event_rows: int, case_rows: int, repeats: int) -> dict[str, Any]: + database = _safe_benchmark_database(dsn) + apply_migrations(dsn) + with psycopg.connect(dsn, autocommit=True) as connection: + _seed_fixture(connection, event_rows, case_rows) + recent_measurements = _profile_one( + connection, + index_name=EVENT_INDEX, + create_index_sql=( + "CREATE INDEX fabops_event_log_event_type_sequence_idx " + "ON fabops_event_log(event_type, sequence DESC)" + ), + query_sql=( + "SELECT sequence, envelope, delivery_status FROM fabops_event_log " + "WHERE event_type = 'process.measurement.recorded.v1' " + "ORDER BY sequence DESC LIMIT %s" + ), + params=(128,), + repeats=repeats, + ) + related_cases = _profile_one( + connection, + index_name=CASE_INDEX, + create_index_sql=( + "CREATE INDEX fabops_cases_classification_lot_updated_idx " + "ON fabops_cases(classification, lot_id DESC, updated_at DESC)" + ), + query_sql=( + "SELECT case_document FROM fabops_cases " + "WHERE classification = %s AND case_id <> %s " + "ORDER BY lot_id DESC, updated_at DESC LIMIT %s" + ), + params=("physical_excursion", "CASE-0000001", 20), + repeats=repeats, + ) + return { + "experiment": "fabops-postgresql-operational-index-profile-v1", + "database": database, + "fixture": {"event_rows": event_rows, "case_rows": case_rows}, + "repeats_per_phase": repeats, + "queries": { + "recent_measurement_events": recent_measurements, + "related_cases": related_cases, + }, + "limitations": [ + "Synthetic portfolio fixture; results are not a production-fab capacity claim.", + "Warm-cache EXPLAIN ANALYZE runs on one local PostgreSQL instance.", + "The benchmark isolates read-plan effects and does not quantify index write amplification.", + ], + } + + +def main() -> None: + parser = argparse.ArgumentParser(description="Profile FabOps PostgreSQL repository query plans before/after selected indexes.") + parser.add_argument("--dsn", required=True) + parser.add_argument("--event-rows", type=int, default=200_000) + parser.add_argument("--case-rows", type=int, default=80_000) + parser.add_argument("--repeats", type=int, default=15) + parser.add_argument("--output") + parser.add_argument("--strict", action="store_true") + args = parser.parse_args() + + result = run_profile( + args.dsn, + event_rows=max(10_000, args.event_rows), + case_rows=max(10_000, args.case_rows), + repeats=max(3, args.repeats), + ) + rendered = json.dumps(result, indent=2, sort_keys=True) + print(rendered) + if args.output: + Path(args.output).write_text(rendered + "\n", encoding="utf-8") + if args.strict: + for query_name, query_result in result["queries"].items(): + expected_index = query_result["index"] + before = query_result["before"] + after = query_result["after"] + if expected_index not in after["index_names"]: + raise SystemExit(f"{query_name}: expected index not used after migration") + if after["p50_ms"] >= before["p50_ms"]: + raise SystemExit(f"{query_name}: p50 did not improve") + if after["shared_hit_blocks"] >= before["shared_hit_blocks"]: + raise SystemExit(f"{query_name}: shared buffer hits did not decrease") + + +if __name__ == "__main__": + main() diff --git a/evidence/postgres/operational-index-profile.json b/evidence/postgres/operational-index-profile.json new file mode 100644 index 0000000..3ad82bd --- /dev/null +++ b/evidence/postgres/operational-index-profile.json @@ -0,0 +1,85 @@ +{ + "database": "fabops_bench", + "experiment": "fabops-postgresql-operational-index-profile-v1", + "fixture": { + "case_rows": 80000, + "event_rows": 200000 + }, + "limitations": [ + "Synthetic portfolio fixture; results are not a production-fab capacity claim.", + "Warm-cache EXPLAIN ANALYZE runs on one local PostgreSQL instance.", + "The benchmark isolates read-plan effects and does not quantify index write amplification." + ], + "queries": { + "recent_measurement_events": { + "after": { + "index_names": [ + "fabops_event_log_event_type_sequence_idx" + ], + "mean_ms": 0.0351, + "p50_ms": 0.033, + "p95_ms": 0.0413, + "plan_nodes": [ + "Limit", + "Index Scan" + ], + "shared_hit_blocks": 14, + "shared_read_blocks": 0, + "variance_ms2": 0.000016 + }, + "before": { + "index_names": [ + "fabops_event_log_pkey" + ], + "mean_ms": 10.2462, + "p50_ms": 10.23, + "p95_ms": 10.627, + "plan_nodes": [ + "Limit", + "Index Scan" + ], + "shared_hit_blocks": 12078, + "shared_read_blocks": 0, + "variance_ms2": 0.058831 + }, + "index": "fabops_event_log_event_type_sequence_idx", + "p50_improvement_percent": 99.68 + }, + "related_cases": { + "after": { + "index_names": [ + "fabops_cases_classification_lot_updated_idx" + ], + "mean_ms": 0.0207, + "p50_ms": 0.018, + "p95_ms": 0.0324, + "plan_nodes": [ + "Limit", + "Index Scan" + ], + "shared_hit_blocks": 46, + "shared_read_blocks": 0, + "variance_ms2": 0.000074 + }, + "before": { + "index_names": [ + "idx_fabops_cases_lot_id" + ], + "mean_ms": 0.1637, + "p50_ms": 0.162, + "p95_ms": 0.1799, + "plan_nodes": [ + "Limit", + "Incremental Sort", + "Index Scan" + ], + "shared_hit_blocks": 2844, + "shared_read_blocks": 0, + "variance_ms2": 0.000095 + }, + "index": "fabops_cases_classification_lot_updated_idx", + "p50_improvement_percent": 88.89 + } + }, + "repeats_per_phase": 15 +} diff --git a/tests/test_postgres_profile.py b/tests/test_postgres_profile.py new file mode 100644 index 0000000..7d73442 --- /dev/null +++ b/tests/test_postgres_profile.py @@ -0,0 +1,37 @@ +from __future__ import annotations + +import pytest + +from evaluation.postgres_profile import _index_names, _percentile, _safe_benchmark_database + + +def test_profile_percentile_is_interpolated_and_deterministic() -> None: + values = [1.0, 2.0, 3.0, 4.0, 100.0] + assert _percentile(values, 0.5) == 3.0 + assert _percentile(values, 0.95) == pytest.approx(80.8) + + +def test_profile_refuses_non_benchmark_database() -> None: + with pytest.raises(SystemExit, match="bench.*test"): + _safe_benchmark_database("postgresql://postgres:postgres@127.0.0.1:5432/fabops") + + +def test_profile_accepts_explicit_benchmark_database() -> None: + assert ( + _safe_benchmark_database("postgresql://postgres:postgres@127.0.0.1:5432/fabops_bench") + == "fabops_bench" + ) + + +def test_profile_collects_nested_index_names() -> None: + plan = { + "Node Type": "Limit", + "Plans": [ + { + "Node Type": "Index Scan", + "Index Name": "fabops_event_log_event_type_sequence_idx", + } + ], + } + + assert _index_names(plan) == ["fabops_event_log_event_type_sequence_idx"]