Agent工作流程入门指南:从零搭建自动化任务处理系统

1次阅读
没有评论

共计 2405 个字符,预计需要花费 7 分钟才能阅读完成。

image.webp

从订单同步看传统轮询的痛点

做过电商系统的同学肯定遇到过这样的场景:用户支付成功后,需要将订单状态同步到物流系统。传统做法往往是这样的:

Agent 工作流程入门指南:从零搭建自动化任务处理系统

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) 模式有三大致命伤:

  1. 资源浪费:90% 的查询都是无效的,数据库压力大
  2. 延迟高:最长要等 1 分钟才能发现状态变更
  3. 扩展难:多个实例同时轮询会导致重复处理

事件驱动架构的优势

改用事件驱动架构 (Event-Driven Architecture) 后,流程变为:

  1. 支付系统完成扣款后发出 PaymentSuccess 事件
  2. 物流 Agent 订阅该事件并触发同步
  3. 库存系统等其他消费者也可以并行处理

这就像餐厅的点餐系统——后厨不需要每隔 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()

关键点:

  1. 锁的 key 要包含业务 ID(如订单号)
  2. 设置合理的超时时间(避免死锁)
  3. 必须先检查是否已处理(网络延迟可能导致锁失效)

死信队列配置

在 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

当需要横向扩展时,注意:

  1. Kafka 按分区 (partition) 数决定并行度
  2. RabbitMQ 使用 consistent_hash 交换器均衡负载
  3. 数据库连接池要随 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)

思考题延伸

如果要实现北京和上海两个数据中心的订单状态最终一致性,可以考虑:

  1. 全局事件 ID(雪花算法 Snowflake)
  2. 跨机房消息同步(如 Kafka MirrorMaker)
  3. 定期对账补偿机制

Agent 系统的魅力在于将复杂的业务流程拆解为可观测、可重试的独立单元。就像乐高积木,每个 Worker 只专注于自己的那块拼图,组合起来却能构建出强大的自动化王国。

正文完
 0
评论(没有评论)