"""
第 17 章 - Kafka Streams 与 ksqlDB
realtime_topn.py - 用 Faust 实现「实时 Top-N 商品销量排行榜」

业务场景：
    电商大屏需要「最近 1 小时下单数 Top-10 商品」。每个订单是一个事件，
    发到 `learn.orders` Topic；Streams 应用维护一个滑动窗口聚合 + Top-N
    Heap，定时把 Top-10 推到 `learn.topn` Topic 给前端 BFF 拉取。

关键技巧：
    1) 1 小时滚动窗口聚合每商品销量（KTable<product_id, count>）
    2) 维护一个全局 Top-N Heap（Table，单 key='topn'）
    3) 每秒输出一次最新 Top-10 到下游 Topic / Web 接口

依赖：
    pip install faust-streaming python-rocksdb

准备：
    kafka-topics.sh ... --create --topic learn.orders --partitions 6 --replication-factor 1
    kafka-topics.sh ... --create --topic learn.topn   --partitions 1 --replication-factor 1

运行：
    faust -A realtime_topn worker -l info --web-port 6068
    curl http://localhost:6068/topn/now
"""
from __future__ import annotations

import faust
import heapq
from datetime import timedelta
from typing import List, Tuple

WINDOW_SIZE = timedelta(hours=1)
EMIT_EVERY  = timedelta(seconds=5)
TOP_N       = 10


class Order(faust.Record, serializer="json"):
    order_id:   str
    product_id: str
    user_id:    str
    amount:     float
    ts:         float    # event time, seconds


class TopNEntry(faust.Record, serializer="json"):
    product_id: str
    count:      int


app = faust.App(
    "realtime-topn",
    broker="kafka://localhost:9092",
    store="rocksdb://",
    table_standby_replicas=1,
)

orders_topic = app.topic("learn.orders", value_type=Order)
topn_topic   = app.topic("learn.topn",   value_type=List[TopNEntry])

# 1 小时滚动窗口聚合：每个 product_id 的下单笔数
hourly_count = (
    app.Table("hourly-count", default=int, partitions=6)
       .tumbling(WINDOW_SIZE, expires=timedelta(hours=2))
       .relative_to_field(Order.ts)
)

# 用一个 single-key Table 存当前 Top-N 快照（方便 IQ 查询）
topn_snapshot = app.Table("topn-snapshot", default=list, partitions=1)


# -------------------------------------------------------------------
# Agent 1：消费订单流，更新窗口聚合
# -------------------------------------------------------------------
@app.agent(orders_topic)
async def aggregate(stream):
    """每来一个订单：对应 product_id 的窗口计数 +1。"""
    async for order in stream:
        # Faust 的 windowed Table：直接用 [] 操作即可
        hourly_count[order.product_id] += 1


# -------------------------------------------------------------------
# Timer：定期扫描全表，算出 Top-N，写到 snapshot + Kafka
# -------------------------------------------------------------------
@app.timer(interval=EMIT_EVERY.total_seconds())
async def emit_topn():
    """周期性把当前窗口的 Top-N 推送到下游 + 内存快照。"""
    # 当前窗口下的所有 (product, count)
    # 注意：Faust 的 windowed Table 遍历需要 .items() + .current()
    pairs: List[Tuple[str, int]] = []
    for k, w in hourly_count.items():
        try:
            cnt = w.current()
            if cnt > 0:
                pairs.append((k, cnt))
        except Exception:
            continue

    # heapq.nlargest = Top-N
    top = heapq.nlargest(TOP_N, pairs, key=lambda x: x[1])
    payload = [TopNEntry(product_id=p, count=c) for p, c in top]

    # 写到 snapshot Table（IQ 查询）
    topn_snapshot["current"] = [(e.product_id, e.count) for e in payload]
    # 写到下游 Topic
    await topn_topic.send(value=payload)

    if payload:
        print(f"[Top-{TOP_N}] " +
              " | ".join(f"{e.product_id}={e.count}" for e in payload))


# -------------------------------------------------------------------
# Web 接口：实时返回 Top-N（Interactive Query）
# -------------------------------------------------------------------
@app.page("/topn/now")
async def http_topn(web, request):
    items = topn_snapshot.get("current", [])
    return web.json([{"product_id": p, "count": c} for p, c in items])


if __name__ == "__main__":
    app.main()


# ============================================================================
# 性能与扩展
# ----------------------------------------------------------------------------
# - 商品总数 1M+ 时，timer 里全表扫描会变慢，建议改成「维护 Top-N 增量更新」
#   （比如只在 cnt 变化时检查是否能替换 heap 末尾）。
# - 窗口越长，State Store 越大；可以加 `compacting` 后端 + 缩短 expires。
# - 真正高吞吐（百万 QPS）建议用 Java Streams 或 Flink；Faust 里 Python
#   单进程吞吐有 GIL 上限。
# ============================================================================
