Skip to content

第 21 章 综合实战项目:实时订单总线 + CDC 数仓链路

收官章。把前 20 章的所有 Kafka 知识串成一套真实可跑的系统

学完你会:独立设计一个中等规模业务的 Kafka 链路,从 Topic 选型、可靠性、Schema 治理、可观测,到 CDC + Streams + Sink 全套落地。

项目特点:

  • 完整可跑docker compose up 即可起一套 Kafka + Schema Registry + Connect + MySQL + ClickHouse
  • 覆盖最多技术点:Producer 幂等 + 多消费组 + Streams 清洗 + Connect Source/Sink + Schema 演进 + 监控
  • 真实业务场景:电商订单 + CDC 数仓双链路

0. 为什么选这个项目?

候选方案难在哪Kafka 能力覆盖
秒杀削峰业务规则多、需要 Redis 协同Producer + 消费组 + 限流
用户行为埋点数据量大但场景单一Producer + Streams
IoT 上报端侧复杂、Kafka 部分较薄多分区 + 路由
订单总线 + CDC 数仓业务+数仓双链路几乎所有:Producer 幂等 / 多消费组 / EOS / Connect / Streams / Schema Registry / 监控

最终选「订单总线 + CDC 数仓」,因为它的每一个子功能正好能映射到前面章节的某个核心知识点——「一个项目讲完所有 Kafka 核心特性」


1. 项目目标与功能清单

1.1 业务背景

某电商订单系统,目标:

  1. 订单事件实时分发:订单创建 / 支付 / 取消 等事件,实时推送给「风控 / 库存 / 通知」三个下游
  2. MySQL CDC 同步:订单表通过 Debezium 捕获变更,经 Kafka Streams 清洗,落到 ClickHouse 做实时分析

1.2 用户故事

角色故事
用户下单 → 立刻收到通知
风控收到订单事件 → 实时识别可疑订单
库存收到订单事件 → 扣库存
数据分析师在 ClickHouse 看到「过去 1 分钟新增订单数」「按城市分组的 GMV」
运维Kafka 集群指标可见、Lag 可监控、有故障演练剧本

1.3 REST 接口清单

方法路径说明
POST/orders创建订单(同步入库 → 异步发 Kafka)
GET/orders/{id}查订单
GET/health健康检查

2. 系统架构

2.1 完整架构图(Mermaid)

2.2 数据流

链路 A:业务事件总线

用户下单 → Order Service
        ├── 1. 同步写 MySQL orders 表(事务)
        └── 2. 异步发 Kafka orders.events(Avro 序列化)
                ├── Risk Consumer    : 实时风控
                ├── Inventory Consumer : 扣库存
                └── Notify Consumer  : 通知用户/商家

链路 B:CDC 数仓

MySQL orders 表
   ↓ binlog
Debezium MySQL Source Connector
   ↓ 写
Kafka topic: cdc.mysql.orders (Avro,schema = before/after/op/source)
   ↓ 消费
Kafka Streams (faust-streaming)
   ↓ 清洗:抽取 after,加入 ip → 城市,统一时区
Kafka topic: orders.enriched
   ↓ 消费
ClickHouse Sink Connector
   ↓ 写
ClickHouse: orders_dwd 表

3. Topic 设计

3.1 核心 Topic

Topic分区副本保留Cleanup说明
orders.events1237ddelete主业务总线,分区数预估 3 年峰值
cdc.mysql.orders3330ddeleteCDC raw,保留时间长方便重放
orders.enriched3314ddelete清洗后用于落仓
orders.events.dlq1330ddelete死信队列
_schemas13compactSchema Registry 内部

3.2 选型理由

分区数 = 12

预估:

  • 当前订单峰值 5K msg/s
  • 单分区合理消费 ≈ 1K msg/s(取决于业务复杂度)
  • → 至少 5 分区
  • 留 2-3 倍余量(业务增长 + 加消费者并发)→ 12

📌 教训呼应案例 1:分区一旦定下来很难加(加分区会破坏 Key 顺序)。宁多勿少,但单 Broker 不超过 2000 分区。

Key 选型 = order_id

  • order_id 是 64 位雪花 ID,高基数
  • 同一订单的所有事件(创建、支付、取消、退款)落同一分区,保证顺序
  • 避免案例 2 的「日期 / 城市」热点

副本数 = 3 + min.insync.replicas = 2

  • 3 副本能容忍 1 台机器宕
  • min.insync.replicas=2 + acks=all = 数据不丢的标配(呼应案例 3)

保留策略

  • 业务 Topic 7 天:覆盖一周回溯需求
  • CDC raw 30 天:方便数仓重放历史
  • DLQ 30 天:足够时间分析死信原因

3.3 Topic 创建命令

init.sh。关键:

bash
kafka-topics.sh --create --bootstrap-server kafka:9092 \
  --topic orders.events --partitions 12 --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config retention.ms=604800000 \
  --config compression.type=zstd \
  --config max.message.bytes=1048576

4. Schema 设计

4.1 orders.events Schema (Avro v1)

json
{
  "type": "record",
  "namespace": "com.shop.events",
  "name": "OrderEvent",
  "fields": [
    {"name": "order_id",   "type": "string"},
    {"name": "user_id",    "type": "long"},
    {"name": "amount",     "type": "double"},
    {"name": "currency",   "type": "string", "default": "CNY"},
    {"name": "status",     "type": {"type":"enum","name":"OrderStatus","symbols":["CREATED","PAID","CANCELED","REFUNDED"]}},
    {"name": "city",       "type": ["null","string"], "default": null},
    {"name": "created_at", "type": {"type":"long","logicalType":"timestamp-millis"}},
    {"name": "items", "type": {"type":"array","items":{
        "type":"record","name":"Item",
        "fields":[
          {"name":"sku","type":"string"},
          {"name":"qty","type":"int"},
          {"name":"price","type":"double"}
        ]}}}
  ]
}

兼容性策略:BACKWARD_TRANSITIVE(呼应第 18 章)

4.2 后续 v2 演进示例

coupon_code 字段,带默认值:

json
{"name": "coupon_code", "type": ["null","string"], "default": null}

→ BACKWARD 通过 ✓


5. 可靠性方案

5.1 Producer 端(Order Service)

python
producer = Producer({
    "bootstrap.servers": "...",
    "acks": "all",                              # 等所有 ISR
    "enable.idempotence": True,                 # PID 去重
    "max.in.flight.requests.per.connection": 5, # 与幂等兼容的最大并发
    "compression.type": "zstd",                 # 压缩省带宽
    "linger.ms": 10,                            # 攒批
    "retries": 2147483647,                      # 实际靠 delivery.timeout.ms 控制
    "delivery.timeout.ms": 120000,              # 2 分钟内必须成功或失败
})

事务模式(如果业务真需要 EOS):可用 transactional.id,但本项目不开(呼应案例 11,只在「绝不能重复」时才用)。

5.2 业务一致性:本地表 + 异步发 Kafka

经典 Outbox Pattern

python
def create_order(order):
    with db.transaction():
        db.execute("INSERT INTO orders ...")           # 1. 入库
        db.execute("INSERT INTO outbox VALUES (...)")  # 2. 写出箱(同事务)
    # 事务提交后,由 Outbox Relay 异步推到 Kafka

📌 不要做的事:在事务外直接 producer.send()。一旦数据库 commit 后 Kafka 发失败,业务和事件不一致。

5.3 Consumer 端

python
consumer = Consumer({
    "bootstrap.servers": "...",
    "group.id": "risk-group",
    "enable.auto.commit": False,             # 手动提交
    "auto.offset.reset": "earliest",         # 新组从头消费历史
    "max.poll.interval.ms": 600000,          # 10 分钟(呼应案例 5)
    "max.poll.records": 100,
    "partition.assignment.strategy": "cooperative-sticky",
})

消费三步骤

  1. poll() 拿一批
  2. 业务幂等处理(DB 唯一约束 / Redis SETNX)
  3. 处理成功后 commit()
python
for msg in batch:
    try:
        process(msg)               # 业务幂等
    except Exception as e:
        send_to_dlq(msg)           # 致命错误进 DLQ
consumer.commit(asynchronous=False)

5.4 死信队列 (DLQ)

任何处理失败超过 N 次的消息进 orders.events.dlq,附带:

  • 原 topic / partition / offset
  • 错误堆栈
  • 处理时间戳

DLQ 由人工或 Job 定期处理。


6. 监控方案

6.1 关键指标 + 阈值

指标阈值告警级别
OfflinePartitionsCount> 0 持续 1 分钟P0
UnderReplicatedPartitions> 0 持续 5 分钟P1
ActiveControllerCount≠ 1P0
UncleanLeaderElectionsPerSec> 0P0
kafka_consumergroup_lag{group="risk-group"}> 5000 持续 5 分钟P1
kafka_consumergroup_lag{group="cdc-streams"}> 50000 持续 10 分钟P2
Order Service Producer error rate> 0.1%P1
Connect status!= RUNNINGP1

6.2 Dashboard

  • Kafka Cluster 总览(JMX Exporter,Grafana ID 11962)
  • Consumer Groups Lag(kafka-exporter,Grafana ID 7589)
  • Connect Workers(自定义看板)
  • 业务指标:订单 TPS / 各下游处理 TPS / 端到端延迟

6.3 日志聚合

所有服务日志通过 Filebeat → ES,关键字告警:

  • ERROR / Exception
  • Rebalance triggered
  • NotEnoughReplicasException

7. 性能基准

7.1 测试环境

  • 3 Broker,每台 4C / 16G / 200G NVMe SSD
  • 1 Schema Registry,1 Connect Worker
  • 1 个 Order Service 实例(FastAPI + uvicorn 4 worker)

7.2 目标

指标目标
Order Service POST /orders QPS5K
Producer p99 延迟< 50 ms
Kafka end-to-end 延迟 (produce → consume)< 100 ms p99
Consumer 组消费速率5K msg/s/组
CDC 链路端到端延迟 (MySQL → ClickHouse)< 10 s p99

7.3 压测方法

locustwrk 压 Order Service:

bash
wrk -t12 -c400 -d60s -s post.lua http://localhost:8080/orders

观察 Kafka JMX:MessagesInPerSec 应稳定在 5K,request-latency-p99 < 50 ms。


8. 故障演练剧本

每季度演练一次,确认系统鲁棒性。

8.1 Broker 宕机

bash
docker stop kafka-2

期望

  • 30 秒内 Leader 切走
  • Order Service 短暂报 NotLeaderForPartition,自动重试成功
  • Lag 短暂上涨后回落

8.2 消费滞后(Lag 暴涨)

人为暂停 Risk Consumer:

bash
docker pause risk-consumer

期望

  • 10 分钟内 P1 告警触发
  • 恢复后 Consumer 能从 committed offset 续上,不丢不重

8.3 Broker 磁盘满

bash
docker exec kafka-1 dd if=/dev/zero of=/var/lib/kafka/data/fill bs=1M count=180000

期望

  • 该 Broker 拒收新写入但不 panic
  • Producer 切走到其他 Broker
  • 清理后能恢复

8.4 网络分区

bash
docker exec kafka-2 iptables -A INPUT -s kafka-1 -j DROP

期望

  • ISR 收缩,UnderReplicated 告警
  • 业务降级(acks=all 阻塞少量分区写入)但不崩溃
  • 恢复后 ISR 自动扩张

8.5 Connect Worker 挂

bash
docker stop connect

期望

  • Debezium 与 ClickHouse Sink 都停止
  • 重启后能从 connect-offsets 续上,不丢 / 不重(Source 幂等)

9. 分阶段迭代路线图

阶段范围关键变化时间
MVPOrder Service + Risk Consumer同步生产 + 单消费组Week 1-2
高可靠+ idempotence + ACL + 监控acks=all + 监控告警上线Week 3-4
多消费组+ Inventory + Notify一份消息多家消费Week 5
Schema 治理接入 Schema RegistryAvro + BACKWARD_TRANSITIVEWeek 6
CDC 数仓+ Debezium + Streams + CK Sink完整 CDC 链路Week 7-9
多活+ MirrorMaker 2 跨机房北京 ↔ 上海 双活Week 10-12

10. 启动手册

10.1 一键启动

bash
cd 21_project
docker compose -f ../docker-compose.yml -f docker-compose-extra.yml up -d
./init.sh                       # 创建 topic + 初始化 MySQL
pip install -r ../requirements.txt

10.2 启动各服务

bash
# Terminal 1: Order Service
python code/order_producer.py

# Terminal 2-4: Consumers
python code/risk_consumer.py
python code/inventory_consumer.py
python code/notify_consumer.py

# Terminal 5: CDC Streams
python code/cdc_streams.py worker

# Register Connectors:
curl -X POST -H "Content-Type: application/json" \
  --data @code/debezium_config.json http://localhost:8083/connectors

curl -X POST -H "Content-Type: application/json" \
  --data @code/clickhouse_sink_config.json http://localhost:8083/connectors

10.3 验证

bash
# 下一笔订单
curl -X POST localhost:8080/orders -H "Content-Type: application/json" \
  -d '{"user_id":1001,"items":[{"sku":"SKU-A","qty":1,"price":99.0}]}'

# 看 ClickHouse 是否落表
docker exec -it clickhouse clickhouse-client \
  -q "SELECT * FROM orders_dwd ORDER BY created_at DESC LIMIT 5"

11. 本章面试高频题

Q1:为什么订单 Topic 选 12 分区?

答案

  • 当前峰值 5K msg/s,单分区合理消费 ≈ 1K msg/s → 至少 5 分区
  • 留 2-3 倍余量给「业务增长 + 提高消费并发」
  • 12 = 5 × 2.4,且 12 容易被 1/2/3/4/6 整除(支持 1-12 个消费者实例都均匀)
  • 不超过单 Broker 推荐上限(< 2000 分区)

加分:能解释「分区数定下来后再扩很难,会破坏 Key 顺序」、能讲清楚分区数与 Broker 数的关系。


Q2:Order Service 同时入 MySQL 和 Kafka,怎么保证一致性?

答案

  • 错误做法:在 DB 事务外 producer.send(),事务 commit 后发失败 → 不一致
  • 正确做法 Outbox Pattern
    1. 把 Kafka 消息内容先写到 outbox 表(与业务表同事务)
    2. 单独的 relay 进程从 outbox 表读取,发到 Kafka
    3. 发送成功后删 outbox 行
  • 更进一步:用 Debezium 直接监听 outbox 表,自动推到 Kafka(这就是 Debezium Outbox Event Router)

加分:能对比「2PC / Saga / Outbox」三种方案,并指出 Outbox 在 Kafka 场景下最合适。


Q3:CDC 链路的 Schema Evolution 怎么管?

答案

  • 全局策略 BACKWARD_TRANSITIVE(任何新版本能解析所有历史版本)
  • 上游 DDL 规约:
    • 加列必须有 DEFAULTNULL
    • 改类型禁止
    • 删列要先在表里 deprecate 30 天
  • Debezium 检测到 DDL 后自动注册新 Schema 版本
  • 下游 Streams 用「最新 Schema」反序列化,旧消息靠默认值填充
  • 不兼容变更走「评审」流程,临时改 NONE 是逃生口

Q4:风控、库存、通知都消费同一个 Topic,需要注意什么?

答案

  • 不同 group.id,每组独立消费完整数据(多对多消费模式)
  • 每组独立 Lag 监控
  • 处理速度差异大时,慢的组可能 Lag 高 → 单独扩容
  • 共享消费(同 group.id 多实例)只在「想分担同一职责」时用
  • 失败处理独立:每组有自己的 DLQ 或重试逻辑
  • 顺序保证只在分区内:同 order_id 的所有事件按顺序,跨 order 不保证

Q5:如何把这套链路升级成「跨机房双活」?

答案

  • MirrorMaker 2 (MM2) 在两个集群间双向复制:
    • 北京 → 上海:北京的 orders.events 在上海集群叫 bj.orders.events
    • 上海 → 北京:同理 sh.orders.events
  • Order Service 写本地集群(acks=all 本地副本),降低跨机房延迟
  • Consumer 优先消费本地集群 + 必要时跨集群消费 mirror 副本
  • 关键问题:避免循环复制(MM2 用 source-cluster prefix 解决)
  • 切换演练:把流量切到上海集群,验证 Lag 在可接受范围

12. 小结

  • 一个真正能跑的项目比 100 个「Hello World」更能教会你 Kafka
  • 核心能力清单(自检表):
    • [ ] Topic 设计:分区数 / 副本 / Key / 保留 / Cleanup 都说得出 why
    • [ ] Producer:acks=all + idempotence + 压缩 + 攒批
    • [ ] Outbox Pattern:保证业务-事件一致性
    • [ ] Consumer:手动提交 + 业务幂等 + DLQ + Cooperative
    • [ ] Schema Registry:BACKWARD_TRANSITIVE + 加字段带默认
    • [ ] CDC:Debezium + Streams + Sink Connector 链路通
    • [ ] 监控:5 大致命指标 + Lag + 业务 TPS
    • [ ] 故障演练:Broker 宕 / Lag 涨 / 磁盘满 / 网络分区

学到这里,你应该有底气说:「Kafka 我能用、能调、能排障、能讲原理」。


🎉 恭喜你完成了从 0 到 1 的 Kafka 之旅。下一步:把每章的代码都跑一遍,把面试题都默写一遍,把项目搭起来。看 100 遍不如动手 1 遍。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
CDC Streams:用 faust-streaming 把 cdc.mysql.orders 清洗后写到 orders.enriched

清洗规则:
1. 抽 after 字段(Debezium envelope 的精华)
2. 过滤 op == 'd'(删除)
3. 加 city:根据 user_id 简单 mock;生产应查 Redis 维表
4. 时区统一 UTC+8

启动:
    python cdc_streams.py worker  -l info

依赖:pip install faust-streaming  (注意:用 fork 而非已停更的 faust)
"""

import os
import json
import faust

BOOTSTRAP = os.getenv("BOOTSTRAP", "kafka://localhost:9092")
SRC_TOPIC = "cdc.mysql.orders"
DST_TOPIC = "orders.enriched"

app = faust.App(
    "cdc-streams",
    broker=BOOTSTRAP,
    value_serializer="json",     # 简化示例:用 JSON。生产改 Avro
    consumer_max_fetch_size=10485760,
    topic_replication_factor=3,
    processing_guarantee="at_least_once",
    store="memory://",
)

src = app.topic(SRC_TOPIC, value_type=bytes)
dst = app.topic(DST_TOPIC)


CITY_BY_USER = {  # mock 维表
    1001: "beijing", 1002: "shanghai", 1003: "shenzhen",
    1004: "hangzhou", 1005: "chengdu",
}


@app.agent(src)
async def process(stream):
    async for raw in stream:
        try:
            env = json.loads(raw)
        except Exception as e:
            print(f"[streams][skip] decode fail: {e}")
            continue

        op = env.get("op") or env.get("payload", {}).get("op")
        after = env.get("after") or env.get("payload", {}).get("after")
        if op == "d" or after is None:
            continue

        # enrichment
        user_id = after.get("user_id")
        if "city" not in after or after["city"] is None:
            after["city"] = CITY_BY_USER.get(user_id, "unknown")

        # 时区:MySQL 默认 UTC,转成 +08:00(这里只是示例)
        # after["created_at"] 假设是字符串
        out = {
            "order_id":   after["order_id"],
            "user_id":    after["user_id"],
            "amount":     float(after["amount"]),
            "currency":   after.get("currency", "CNY"),
            "status":     after.get("status", "CREATED"),
            "city":       after["city"],
            "items_json": after.get("items_json"),
            "created_at": after.get("created_at"),
            "updated_at": after.get("updated_at"),
        }
        await dst.send(value=out)


if __name__ == "__main__":
    app.main()
python
#!/usr/bin/env python3
"""
Inventory Consumer:库存消费者

特点:
- 业务幂等:用 (order_id, sku) 作为 DB 唯一约束,重复消费不会扣两次
- 失败 ≥ 3 次 → DLQ
"""

import os
import time
import json
import pymysql
from confluent_kafka import Consumer, Producer
from confluent_kafka.serialization import SerializationContext, MessageField
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"
DLQ = "orders.events.dlq"
GROUP = "inventory-group"

MY_CFG = dict(
    host=os.getenv("MY_HOST", "127.0.0.1"),
    port=int(os.getenv("MY_PORT", "3306")),
    user=os.getenv("MY_USER", "root"),
    password=os.getenv("MY_PWD", "root"),
    database="shop",
    autocommit=True,
)


def ensure_dedup_table():
    conn = pymysql.connect(**MY_CFG)
    with conn.cursor() as cur:
        cur.execute("""
            CREATE TABLE IF NOT EXISTS inventory_dedup (
              order_id VARCHAR(64),
              sku      VARCHAR(64),
              qty      INT,
              applied_at DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3),
              PRIMARY KEY (order_id, sku)
            ) ENGINE=InnoDB
        """)
    conn.close()


def apply_inventory(event):
    """利用 PRIMARY KEY 实现幂等。"""
    conn = pymysql.connect(**MY_CFG)
    try:
        with conn.cursor() as cur:
            for item in event.get("items", []):
                try:
                    cur.execute(
                        "INSERT INTO inventory_dedup (order_id, sku, qty) VALUES (%s, %s, %s)",
                        (event["order_id"], item["sku"], item["qty"]),
                    )
                    # 这里可以接真实库存表的扣减
                    print(f"[inv]   - 扣 sku={item['sku']} qty={item['qty']}")
                except pymysql.IntegrityError:
                    print(f"[inv]   = 已扣过,幂等跳过 sku={item['sku']}")
    finally:
        conn.close()


def main():
    ensure_dedup_table()

    sr = SchemaRegistryClient({"url": SR_URL})
    deser = AvroDeserializer(sr, from_dict=lambda d, ctx: d)
    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": 600_000,
        "max.poll.records": 50,
        "partition.assignment.strategy": "cooperative-sticky",
    })
    dlq_p = Producer({"bootstrap.servers": BOOTSTRAP})
    consumer.subscribe([TOPIC])
    print(f"[inv] start, group={GROUP}")

    retry_count = {}
    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                continue

            key = (msg.partition(), msg.offset())
            try:
                event = deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))
                if event["status"] != "CREATED":
                    consumer.commit(msg, asynchronous=False)
                    continue

                print(f"[inv] order={event['order_id']} 扣库存 ...")
                apply_inventory(event)
                consumer.commit(msg, asynchronous=False)
                retry_count.pop(key, None)
            except Exception as e:
                retry_count[key] = retry_count.get(key, 0) + 1
                print(f"[inv][ERR] retry {retry_count[key]}/3: {e}")
                if retry_count[key] >= 3:
                    print(f"[inv] → DLQ {key}")
                    dlq_p.produce(DLQ, key=msg.key(), value=msg.value(),
                                  headers=[("error", str(e).encode()),
                                           ("origin", b"inventory")])
                    dlq_p.flush(5)
                    consumer.commit(msg, asynchronous=False)
                    retry_count.pop(key, None)
                else:
                    time.sleep(1)
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
Notify Consumer:通知消费者
- 模拟向用户/商家发短信/IM 通知
- 失败有限次重试 → DLQ
"""

import os
import random
import time
from confluent_kafka import Consumer, Producer
from confluent_kafka.serialization import SerializationContext, MessageField
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"
DLQ = "orders.events.dlq"
GROUP = "notify-group"


def send_notification(event):
    """模拟通知,5% 概率失败。"""
    if random.random() < 0.05:
        raise RuntimeError("notification gateway timeout")
    print(f"[notify] 📧 user={event['user_id']} order={event['order_id']} 已通知")


def main():
    sr = SchemaRegistryClient({"url": SR_URL})
    deser = AvroDeserializer(sr, from_dict=lambda d, ctx: d)
    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": 300_000,
        "partition.assignment.strategy": "cooperative-sticky",
    })
    dlq_p = Producer({"bootstrap.servers": BOOTSTRAP})
    consumer.subscribe([TOPIC])
    print(f"[notify] start, group={GROUP}")

    retries = {}
    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None or msg.error():
                continue

            key = (msg.partition(), msg.offset())
            try:
                event = deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))
                send_notification(event)
                consumer.commit(msg, asynchronous=False)
                retries.pop(key, None)
            except Exception as e:
                retries[key] = retries.get(key, 0) + 1
                if retries[key] < 3:
                    print(f"[notify][ERR] retry {retries[key]}/3: {e}")
                    time.sleep(0.5)
                else:
                    print(f"[notify] → DLQ after 3 retries")
                    dlq_p.produce(DLQ, key=msg.key(), value=msg.value(),
                                  headers=[("error", str(e).encode()),
                                           ("origin", b"notify")])
                    dlq_p.flush(5)
                    consumer.commit(msg, asynchronous=False)
                    retries.pop(key, None)
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
Order Service:FastAPI + Kafka Producer

职责:
1. 接收 HTTP POST /orders
2. 同步入 MySQL(事务)
3. 异步发 Kafka orders.events(Avro,带 idempotence)

启动:
    uvicorn order_producer:app --host 0.0.0.0 --port 8080 --workers 4

依赖:
    pip install fastapi uvicorn pymysql confluent-kafka[avro]
"""

import os
import json
import time
import uuid
from contextlib import asynccontextmanager

import pymysql
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field
from typing import List, Optional

from confluent_kafka import Producer
from confluent_kafka.serialization import SerializationContext, MessageField, StringSerializer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"

MY_HOST = os.getenv("MY_HOST", "127.0.0.1")
MY_PORT = int(os.getenv("MY_PORT", "3306"))
MY_USER = os.getenv("MY_USER", "root")
MY_PWD = os.getenv("MY_PWD", "root")
MY_DB = "shop"

ORDER_EVENT_SCHEMA = """
{
  "type": "record",
  "namespace": "com.shop.events",
  "name": "OrderEvent",
  "fields": [
    {"name": "order_id",   "type": "string"},
    {"name": "user_id",    "type": "long"},
    {"name": "amount",     "type": "double"},
    {"name": "currency",   "type": "string", "default": "CNY"},
    {"name": "status",     "type": {"type":"enum","name":"OrderStatus",
              "symbols":["CREATED","PAID","CANCELED","REFUNDED"]}},
    {"name": "city",       "type": ["null","string"], "default": null},
    {"name": "created_at", "type": {"type":"long","logicalType":"timestamp-millis"}},
    {"name": "items", "type": {"type":"array","items":{
        "type":"record","name":"Item",
        "fields":[
          {"name":"sku","type":"string"},
          {"name":"qty","type":"int"},
          {"name":"price","type":"double"}
        ]}}}
  ]
}
"""


class Item(BaseModel):
    sku: str
    qty: int = Field(ge=1)
    price: float = Field(ge=0)


class CreateOrder(BaseModel):
    user_id: int
    items: List[Item]
    city: Optional[str] = None


# 全局对象
producer: Producer = None
avro_serializer: AvroSerializer = None
key_serializer = StringSerializer("utf_8")
db_pool = None


def get_conn():
    return pymysql.connect(
        host=MY_HOST, port=MY_PORT, user=MY_USER, password=MY_PWD,
        database=MY_DB, autocommit=False, charset="utf8mb4",
    )


@asynccontextmanager
async def lifespan(app: FastAPI):
    global producer, avro_serializer
    sr = SchemaRegistryClient({"url": SR_URL})
    avro_serializer = AvroSerializer(sr, ORDER_EVENT_SCHEMA, lambda x, ctx: x)
    producer = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "acks": "all",
        "enable.idempotence": True,
        "max.in.flight.requests.per.connection": 5,
        "linger.ms": 10,
        "compression.type": "zstd",
        "delivery.timeout.ms": 120000,
        "client.id": f"order-svc-{uuid.uuid4().hex[:8]}",
    })
    yield
    producer.flush(10)


app = FastAPI(title="Order Service", lifespan=lifespan)


def _delivery_cb(err, msg):
    if err is not None:
        # 实际生产应记录到本地 disk / outbox 重试
        print(f"  [!!!] kafka send failed: {err}")
    else:
        print(f"  [✓] kafka offset = {msg.offset()} partition = {msg.partition()}")


@app.get("/health")
def health():
    return {"status": "ok"}


@app.post("/orders")
def create_order(o: CreateOrder):
    order_id = "ORD-" + uuid.uuid4().hex[:16].upper()
    amount = round(sum(i.qty * i.price for i in o.items), 2)
    items_json = json.dumps([i.model_dump() for i in o.items])

    # 1. 入 MySQL(事务)
    conn = get_conn()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "INSERT INTO orders (order_id, user_id, amount, status, city, items_json) "
                "VALUES (%s, %s, %s, %s, %s, %s)",
                (order_id, o.user_id, amount, "CREATED", o.city, items_json),
            )
            # Outbox 模式(可选):把事件先写出箱表
            cur.execute(
                "INSERT INTO outbox (aggregate_id, event_type, payload) VALUES (%s, %s, %s)",
                (order_id, "OrderCreated", json.dumps({
                    "order_id": order_id, "user_id": o.user_id, "amount": amount,
                    "city": o.city,
                })),
            )
        conn.commit()
    except Exception as e:
        conn.rollback()
        raise HTTPException(500, f"db error: {e}")
    finally:
        conn.close()

    # 2. 异步发 Kafka(这里直接发;生产可由 Outbox Relay 单独进程发)
    event = {
        "order_id": order_id,
        "user_id": o.user_id,
        "amount": amount,
        "currency": "CNY",
        "status": "CREATED",
        "city": o.city,
        "created_at": int(time.time() * 1000),
        "items": [i.model_dump() for i in o.items],
    }
    try:
        value = avro_serializer(event, SerializationContext(TOPIC, MessageField.VALUE))
        producer.produce(
            topic=TOPIC,
            key=key_serializer(order_id, SerializationContext(TOPIC, MessageField.KEY)),
            value=value,
            on_delivery=_delivery_cb,
        )
        producer.poll(0)
    except Exception as e:
        # 发送失败不能让用户感知(已入库);记日志,由 Outbox Relay 兜底
        print(f"[WARN] enqueue kafka fail (will be relayed by outbox): {e}")

    return {"order_id": order_id, "status": "CREATED", "amount": amount}


@app.get("/orders/{order_id}")
def get_order(order_id: str):
    conn = get_conn()
    try:
        with conn.cursor(pymysql.cursors.DictCursor) as cur:
            cur.execute("SELECT * FROM orders WHERE order_id=%s", (order_id,))
            row = cur.fetchone()
            if not row:
                raise HTTPException(404, "not found")
            return row
    finally:
        conn.close()


if __name__ == "__main__":
    import uvicorn
    uvicorn.run("order_producer:app", host="0.0.0.0", port=8080, workers=1, reload=False)
python
#!/usr/bin/env python3
"""
Risk Consumer:风控消费者

业务规则(示例):
- amount > 50000 → 标记可疑
- 同 user_id 1 分钟内 > 10 笔 → 标记可疑

依赖:pip install confluent-kafka[avro]
"""

import os
import time
from collections import defaultdict, deque
from confluent_kafka import Consumer, Producer
from confluent_kafka.serialization import SerializationContext, MessageField
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"
DLQ = "orders.events.dlq"
GROUP = "risk-group"

WINDOW_MS = 60_000
WINDOW_LIMIT = 10
AMOUNT_LIMIT = 50_000


def main():
    sr = SchemaRegistryClient({"url": SR_URL})
    deser = AvroDeserializer(sr, from_dict=lambda d, ctx: d)

    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": 600_000,
        "max.poll.records": 100,
        "session.timeout.ms": 45000,
        "partition.assignment.strategy": "cooperative-sticky",
    })
    dlq_producer = Producer({"bootstrap.servers": BOOTSTRAP})
    consumer.subscribe([TOPIC])

    user_window = defaultdict(lambda: deque())
    print(f"[risk] start, group={GROUP}")

    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                print(f"[risk] err: {msg.error()}")
                continue

            try:
                event = deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))
            except Exception as e:
                print(f"[risk] decode fail → DLQ: {e}")
                dlq_producer.produce(DLQ, key=msg.key(), value=msg.value(),
                                     headers=[("error", str(e).encode())])
                consumer.commit(msg)
                continue

            # ===== 业务规则 =====
            risky = False
            reasons = []
            if event["amount"] > AMOUNT_LIMIT:
                risky = True
                reasons.append(f"amount {event['amount']} > {AMOUNT_LIMIT}")

            uid = event["user_id"]
            now = time.time() * 1000
            dq = user_window[uid]
            while dq and dq[0] < now - WINDOW_MS:
                dq.popleft()
            dq.append(now)
            if len(dq) > WINDOW_LIMIT:
                risky = True
                reasons.append(f"user {uid} {len(dq)} orders/min")

            tag = "🚨 RISKY" if risky else "✓ ok"
            print(f"[risk] {tag} order={event['order_id']} amt={event['amount']} reasons={reasons}")

            # ===== 提交 offset =====
            consumer.commit(msg, asynchronous=False)
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    main()

cdc_streams.py ↗ · inventory_consumer.py ↗ · notify_consumer.py ↗ · order_producer.py ↗ · risk_consumer.py ↗