#!/usr/bin/env python3
"""
第 7 章 · 三种写入模式对比

  1. 同步逐条 INSERT（反模式，仅作教学）
  2. async_insert（服务端攒批）
  3. Buffer 引擎（内存缓冲表）

输出每种模式的耗时、Part 数、写入吞吐。

依赖：pip install clickhouse-connect
"""
from __future__ import annotations

import random
import time

import clickhouse_connect

HOST = "127.0.0.1"
PORT = 8123
USER = "default"
TOTAL = 100_000
URLS = ["/home", "/p/1", "/p/2", "/login"]


def gen_rows(n: int):
    base = int(time.time())
    for i in range(n):
        yield (
            base + (i % 86400),
            random.randint(1, 100_000),
            random.choice(URLS),
            "direct",
            random.randint(0, 5_000),
        )


def parts(client, table: str) -> int:
    sql = (
        "SELECT count() FROM system.parts "
        f"WHERE database='learn_ck' AND table='{table}' AND active"
    )
    return int(client.query(sql).result_rows[0][0])


def truncate(client, table: str) -> None:
    client.command(f"TRUNCATE TABLE learn_ck.{table}")


def mode_sync_one_by_one(client) -> None:
    """同步逐条 INSERT —— 教学反例，行数较小（5000）以免跑半天。"""
    table = "events_perf"
    n = 5000
    truncate(client, table)
    t0 = time.time()
    for row in gen_rows(n):
        client.command(
            "INSERT INTO learn_ck.events_perf (ts, uid, url, referer, duration) "
            "VALUES",
            parameters=row,
        )
    cost = time.time() - t0
    print(f"[1] sync 1-by-1     n={n:>7d}  cost={cost:7.2f}s  "
          f"parts={parts(client, table):>4d}  qps={n/cost:>10.0f}/s")


def mode_async_insert(client) -> None:
    """async_insert：服务端攒批，wait_for_async_insert=0 让客户端立即返回。"""
    table = "events_async"
    truncate(client, table)
    cols = ["ts", "uid", "url", "referer", "duration"]
    settings = {
        "async_insert": 1,
        "wait_for_async_insert": 0,
        "async_insert_busy_timeout_ms": 1000,
    }
    t0 = time.time()
    for row in gen_rows(TOTAL):
        client.insert(f"learn_ck.{table}", [row],
                      column_names=cols, settings=settings)
    cost = time.time() - t0
    time.sleep(2)
    print(f"[2] async_insert    n={TOTAL:>7d}  cost={cost:7.2f}s  "
          f"parts={parts(client, table):>4d}  qps={TOTAL/cost:>10.0f}/s")


def mode_buffer_engine(client) -> None:
    """Buffer 引擎：写 events_buffer 自动刷到 events_perf。"""
    target = "events_perf"
    truncate(client, target)
    cols = ["ts", "uid", "url", "referer", "duration"]
    t0 = time.time()
    rows = list(gen_rows(TOTAL))
    chunk = 5_000
    for i in range(0, len(rows), chunk):
        client.insert("learn_ck.events_buffer", rows[i:i + chunk],
                      column_names=cols)
    cost = time.time() - t0
    client.command("OPTIMIZE TABLE learn_ck.events_buffer")
    time.sleep(1)
    print(f"[3] Buffer engine   n={TOTAL:>7d}  cost={cost:7.2f}s  "
          f"parts(target)={parts(client, target):>4d}  "
          f"qps={TOTAL/cost:>10.0f}/s")


def main() -> None:
    client = clickhouse_connect.get_client(host=HOST, port=PORT, username=USER)
    print(f"[insert_modes] connected to {HOST}:{PORT}")
    print("=" * 80)
    mode_sync_one_by_one(client)
    mode_async_insert(client)
    mode_buffer_engine(client)
    print("=" * 80)
    print("提示：Part 数越少越好；同行数下 batch_size 越大、吞吐越高。")


if __name__ == "__main__":
    main()
