共计 2405 个字符,预计需要花费 7 分钟才能阅读完成。
从订单同步看传统轮询的痛点
做过电商系统的同学肯定遇到过这样的场景:用户支付成功后,需要将订单状态同步到物流系统。传统做法往往是这样的:

while True:
unpaid_orders = Order.query.filter_by(status='unpaid')
for order in unpaid_orders:
if check_payment(order.id):
update_logistics(order.id)
sleep(60) # 每分钟轮询一次
这种轮询 (polling) 模式有三大致命伤:
- 资源浪费:90% 的查询都是无效的,数据库压力大
- 延迟高:最长要等 1 分钟才能发现状态变更
- 扩展难:多个实例同时轮询会导致重复处理
事件驱动架构的优势
改用事件驱动架构 (Event-Driven Architecture) 后,流程变为:
- 支付系统完成扣款后发出
PaymentSuccess事件 - 物流 Agent 订阅该事件并触发同步
- 库存系统等其他消费者也可以并行处理
这就像餐厅的点餐系统——后厨不需要每隔 5 分钟去前台问 ” 有新订单吗?”,而是前台直接按铃通知。
技术选型对比
| 方案 | 吞吐量(消息 / 秒) | 平均延迟 | 适用场景 |
|---|---|---|---|
| Celery+RabbitMQ | 5,000-10,000 | <100ms | 任务队列,复杂路由需求 |
| Kafka+Redis Stream | 100,000+ | <10ms | 高吞吐流处理 |
新手建议 :从 Celery 开始(版本 >=5.2),它的@task 装饰器让基础用法非常简单:
# 生产者示例
@app.task(bind=True, max_retries=3)
def publish_payment_event(self, order_id):
try:
payment = Payment.get(order_id)
event = {'type': 'payment_success', 'data': payment.to_dict()}
rabbitmq.publish(exchange='orders', routing_key='payment', body=event)
except Exception as e:
self.retry(exc=e, countdown=60) # 1 分钟后重试
消费者 Worker 的幂等设计
最怕的就是重复消费导致库存多扣。Redis 分布式锁是经典解决方案:
# 消费者示例(使用 redis-py>=4.3)def handle_payment_event(event):
lock = redis.lock(f"order_lock:{event['order_id']}", timeout=300)
try:
if lock.acquire(blocking=False):
if not check_processed(event['event_id']):
process_order(event)
mark_processed(event['event_id'])
finally:
lock.release()
关键点:
- 锁的 key 要包含业务 ID(如订单号)
- 设置合理的超时时间(避免死锁)
- 必须先检查是否已处理(网络延迟可能导致锁失效)
死信队列配置
在 Celery 中配置死信队列(Dead Letter Queue):
app.conf.task_queues = {
Queue('payment',
exchange=Exchange('orders'),
routing_key='payment',
queue_arguments={
'x-dead-letter-exchange': 'dlx',
'x-dead-letter-routing-key': 'payment.failed'
})
}
这样处理失败的消息会自动转到 dlx 交换机,方便后续排查。
性能测试与优化
用 Locust(版本 >=2.12)测试单 Worker 吞吐量:
from locust import HttpUser, task
class AgentUser(HttpUser):
@task
def post_event(self):
self.client.post("/events", json={
"type": "test",
"data": {"id": "123"}
})
启动命令:
locust -f test.py --headless -u 1000 -r 100 --run-time 1h
当需要横向扩展时,注意:
- Kafka 按分区 (partition) 数决定并行度
- RabbitMQ 使用
consistent_hash交换器均衡负载 - 数据库连接池要随 Worker 数调整
避坑指南
消息积压告警
在 Prometheus 中设置告警规则:
- alert: MessageBacklog
expr: rate(rabbitmq_queue_messages_unacked[5m]) > 1000
for: 10m
labels:
severity: critical
崩溃恢复
定期保存消费位移(offset):
# 使用 Redis 记录 checkpoint
def save_checkpoint(consumer_id, offset):
redis.hset('agent:checkpoints', consumer_id, offset)
def recover_consumer(consumer_id):
last_offset = redis.hget('agent:checkpoints', consumer_id) or 0
consumer.seek_to(last_offset)
思考题延伸
如果要实现北京和上海两个数据中心的订单状态最终一致性,可以考虑:
- 全局事件 ID(雪花算法 Snowflake)
- 跨机房消息同步(如 Kafka MirrorMaker)
- 定期对账补偿机制
Agent 系统的魅力在于将复杂的业务流程拆解为可观测、可重试的独立单元。就像乐高积木,每个 Worker 只专注于自己的那块拼图,组合起来却能构建出强大的自动化王国。
正文完
