"""
Ch4 配套代码 2 / 3 —— 物化视图实战

依赖：先 init.sql 建好 ch4_sales / ch4_mv_daily_sales

演示：
  1. 普通视图 vs 物化视图查询耗时
  2. REFRESH MATERIALIZED VIEW (锁) vs REFRESH ... CONCURRENTLY (不锁读)
  3. 没有 UNIQUE 索引时 CONCURRENTLY 会报错
  4. 给物化视图加索引让查询走 Index Scan
"""

import time

import psycopg


DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"


def section(title: str) -> None:
    print("\n" + "=" * 64)
    print(title)
    print("=" * 64)


def time_query(conn: psycopg.Connection, sql: str, repeat: int = 3) -> float:
    """跑 N 次取平均，返回毫秒"""
    times = []
    with conn.cursor() as cur:
        for _ in range(repeat):
            t0 = time.perf_counter()
            cur.execute(sql)
            cur.fetchall()
            times.append((time.perf_counter() - t0) * 1000)
    return sum(times) / len(times)


def demo_view_vs_mv(conn: psycopg.Connection) -> None:
    section("Demo 1: 普通视图 vs 物化视图查询耗时")

    sql_view = "SELECT * FROM ch4_v_sales_by_region ORDER BY total DESC"
    sql_mv   = "SELECT region, SUM(total) AS total FROM ch4_mv_daily_sales GROUP BY region ORDER BY total DESC"

    t1 = time_query(conn, sql_view)
    t2 = time_query(conn, sql_mv)

    print(f"  普通视图   ch4_v_sales_by_region   平均 {t1:.2f} ms (每次都算 GROUP BY)")
    print(f"  物化视图   ch4_mv_daily_sales 聚合 平均 {t2:.2f} ms (查表 + 小聚合)")
    if t1 > 0:
        print(f"  💡 物化视图比普通视图快 {t1/t2:.1f} 倍 (基于 5 万行数据，越多差距越大)")


def demo_refresh_modes(conn: psycopg.Connection) -> None:
    section("Demo 2: REFRESH 两种模式")

    with conn.cursor() as cur:
        # 普通 REFRESH（锁 ACCESS EXCLUSIVE）
        t0 = time.perf_counter()
        cur.execute("REFRESH MATERIALIZED VIEW ch4_mv_daily_sales")
        full_ms = (time.perf_counter() - t0) * 1000
        print(f"  REFRESH MATERIALIZED VIEW (FULL):       {full_ms:.1f} ms (期间 SELECT 阻塞)")

        # CONCURRENTLY（不锁读）
        t0 = time.perf_counter()
        cur.execute("REFRESH MATERIALIZED VIEW CONCURRENTLY ch4_mv_daily_sales")
        concurrent_ms = (time.perf_counter() - t0) * 1000
        print(f"  REFRESH MATERIALIZED VIEW CONCURRENTLY: {concurrent_ms:.1f} ms (期间 SELECT 不阻塞)")
        print(f"  💡 CONCURRENTLY 比 FULL 慢 (要做差异比对)，但生产几乎都用它")


def demo_concurrently_requires_unique(conn: psycopg.Connection) -> None:
    section("Demo 3: CONCURRENTLY 要求 UNIQUE 索引")

    with conn.cursor() as cur:
        # 临时创建一个没 UNIQUE 索引的物化视图
        cur.execute("DROP MATERIALIZED VIEW IF EXISTS ch4_mv_no_unique")
        cur.execute("CREATE MATERIALIZED VIEW ch4_mv_no_unique AS SELECT region, count(*) AS cnt FROM ch4_sales GROUP BY region")

        try:
            cur.execute("REFRESH MATERIALIZED VIEW CONCURRENTLY ch4_mv_no_unique")
            print("  ⚠️ 未预期到通过！")
        except psycopg.errors.InvalidTableDefinition as e:
            conn.rollback()
            print(f"  ✅ 预期失败: {str(e).splitlines()[0]}")
            print("     -> 必须 CREATE UNIQUE INDEX 才能 CONCURRENTLY 刷新")

        cur.execute("DROP MATERIALIZED VIEW ch4_mv_no_unique")
        conn.commit()


def demo_explain(conn: psycopg.Connection) -> None:
    section("Demo 4: 物化视图查询的 EXPLAIN")

    with conn.cursor() as cur:
        cur.execute("""
            EXPLAIN (ANALYZE, BUFFERS)
            SELECT * FROM ch4_mv_daily_sales
            WHERE day = current_date - 1 AND region = 'North'
        """)
        print("  ↓ 物化视图按主键命中 Index Scan")
        for row in cur.fetchall():
            print(f"     {row[0]}")


def demo_data_freshness(conn: psycopg.Connection) -> None:
    section("Demo 5: 物化视图的「过期」演示")

    with conn.cursor() as cur:
        # 插入一行新销售记录
        cur.execute("""
            INSERT INTO ch4_sales (region, amount, sold_at)
            VALUES ('North', 9999.99, now())
            RETURNING id
        """)
        new_id = cur.fetchone()[0]
        conn.commit()

        cur.execute("SELECT count(*) FROM ch4_sales WHERE region = 'North'")
        live = cur.fetchone()[0]
        cur.execute("SELECT SUM(cnt) FROM ch4_mv_daily_sales WHERE region = 'North'")
        mv_cnt = cur.fetchone()[0]

        print(f"  写入新记录后：")
        print(f"     ch4_sales (实时)            North 总条数 = {live}")
        print(f"     ch4_mv_daily_sales (过期)   North 总条数 = {mv_cnt}    ← 还没刷新")

        cur.execute("REFRESH MATERIALIZED VIEW CONCURRENTLY ch4_mv_daily_sales")
        cur.execute("SELECT SUM(cnt) FROM ch4_mv_daily_sales WHERE region = 'North'")
        mv_cnt2 = cur.fetchone()[0]
        print(f"     REFRESH 后                  North 总条数 = {mv_cnt2}    ← 同步上了")

        # 清理
        cur.execute("DELETE FROM ch4_sales WHERE id = %s", (new_id,))
        cur.execute("REFRESH MATERIALIZED VIEW CONCURRENTLY ch4_mv_daily_sales")
        conn.commit()


def main() -> None:
    with psycopg.connect(DSN, autocommit=False) as conn:
        demo_view_vs_mv(conn)
        demo_refresh_modes(conn)
        demo_concurrently_requires_unique(conn)
        demo_explain(conn)
        demo_data_freshness(conn)


if __name__ == "__main__":
    try:
        main()
    except psycopg.OperationalError as e:
        print(f"❌ 连接 PostgreSQL 失败: {e}")
