#!/usr/bin/env python3
"""
Rebalance 风暴复现脚本（案例 5）。

操作：
1. 起 N 个 Consumer，故意把每条消息处理设置为 5 分钟（> max.poll.interval.ms 默认 5 分钟）
2. 你会观察到：每个 Consumer 都在第 5 分钟左右被踢，触发 Rebalance
3. Rebalance 后分区重新分配，新接手的 Consumer 也会因相同原因被踢
4. 整个 group 永远在 Rebalance 状态

可调环境变量：
    PROCESS_SECS=400        # 每条消息处理时间（秒），> 300 必现风暴
    MAX_POLL_MS=300000      # max.poll.interval.ms（毫秒）
    INSTANCE_ID=worker-1    # 多终端起多实例

修复方式（取消注释 fix block）：
- PROCESS_SECS 调小 < MAX_POLL_MS / 1000
- 改用 CooperativeStickyAssignor 减小 STW
- max.poll.records 调小

依赖：pip install confluent-kafka
"""

import os
import time
from confluent_kafka import Consumer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
TOPIC = os.getenv("TOPIC", "learn.20.storm")
GROUP = os.getenv("GROUP", "storm-group")
INSTANCE_ID = os.getenv("INSTANCE_ID", f"worker-{os.getpid()}")
PROCESS_SECS = int(os.getenv("PROCESS_SECS", "400"))
MAX_POLL_MS = int(os.getenv("MAX_POLL_MS", "300000"))


def on_assign(c, parts):
    print(f"  [{INSTANCE_ID}] ASSIGN -> {[(p.topic, p.partition) for p in parts]}")

def on_revoke(c, parts):
    print(f"  [{INSTANCE_ID}] REVOKE <- {[(p.topic, p.partition) for p in parts]}")


def main():
    cfg = {
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": MAX_POLL_MS,
        "session.timeout.ms": 45000,
        # === 复现关键点：单次 poll 处理 PROCESS_SECS 秒 ===
        # === 修复方案 1: 改成 CooperativeStickyAssignor ===
        # "partition.assignment.strategy": "cooperative-sticky",
        # === 修复方案 2: 用 group.instance.id 静态成员，避免短抖动触发 ===
        # "group.instance.id": INSTANCE_ID,
    }
    c = Consumer(cfg)
    c.subscribe([TOPIC], on_assign=on_assign, on_revoke=on_revoke)

    print(f"[{INSTANCE_ID}] start, MAX_POLL={MAX_POLL_MS}ms, PROCESS={PROCESS_SECS}s")
    print(f"[{INSTANCE_ID}] 预期：处理时间 > max.poll.interval.ms 时会被踢出 group")

    while True:
        msg = c.poll(1.0)
        if msg is None or msg.error():
            continue

        print(f"[{INSTANCE_ID}] consume p={msg.partition()} off={msg.offset()} → 模拟处理 {PROCESS_SECS}s")
        # ↓↓↓ 故意慢处理 ↓↓↓
        time.sleep(PROCESS_SECS)

        try:
            c.commit(msg)
            print(f"[{INSTANCE_ID}] commit OK")
        except Exception as e:
            print(f"[{INSTANCE_ID}] commit FAIL: {e}  ← 已经被踢出 group")


if __name__ == "__main__":
    try:
        main()
    except KeyboardInterrupt:
        pass
