影刀RPA 消息队列自动化:RabbitMQ Kafka可靠性保证
📅 2026/7/21 6:37:22
👁️ 阅读次数
📝 编程学习
影刀RPA 消息队列自动化:RabbitMQ Kafka可靠性保证
什么情况用什么 → 怎么做 → 有什么坑
作者:林焱 | 飞行社出品
什么情况用什么
用RPA处理业务时,需要和生产系统的消息队列对接——要么从队列取任务,要么推送任务。但消息丢了、重复消费了,排查起来要命。
这套方案适合:
- RPA流程和消息队列集成(异步任务处理)
- 保证消息不丢、不重复处理
- 监控队列堆积情况,及时处理
核心工具:影刀RPA + pika(RabbitMQ) + kafka-python(Kafka) + 监控告警
怎么做
拼多多店群自动化报活动上架!
第一步:RabbitMQ可靠消费(手动ACK + 重试)
importpikaimportjsonfromdatetimeimportdatetimeimporttimeclassReliableRabbitConsumer:"""RabbitMQ可靠消费者:保证消息不丢失"""def__init__(self,host,queue_name,username='guest',password='guest'):self.host=host self.queue_name=queue_name self.credentials=pika.PlainCredentials(username,password)self.connection=Noneself.channel=Nonedefconnect(self):"""建立连接(带重连机制)"""try:parameters=pika.ConnectionParameters(host=self.host,credentials=self.credentials,heartbeat=600,# 心跳超时10分钟blocked_connection_timeout=300)self.connection=pika.BlockingConnection(parameters)self.channel=self.connection.channel()# 声明队列(幂等,已存在不会报错)self.channel.queue_declare(queue=self.queue_name,durable=True,# 队列持久化arguments={'x-death-letter-exchange':'dlx.exchange',# 死信交换机'x-message-ttl':24*3600*1000# 消息TTL 24小时})print(f"✅ RabbitMQ连接成功:{self.host}")returnTrueexceptExceptionase:print(f"⚠️ RabbitMQ连接失败:{e}")returnFalsedefconsume_with_retry(self,process_func,max_retries=3):""" 可靠消费:手动ACK + 失败重试 process_func: 业务处理函数,返回True表示处理成功 """defon_message(channel,method,properties,body):message_id=properties.message_idormethod.delivery_tag retry_count=properties.headers.get('x-retry-count',0)ifproperties.headerselse0try:# 解析消息ifproperties.content_type=='application/json':data=json.loads(body)else:data=body.decode('utf-8')print(f"收到消息{message_id}:{str(data)[:100]}")# 处理业务success=process_func(data)ifsuccess:# 处理成功,手动ACKchannel.basic_ack(delivery_tag=method.delivery_tag)print(f"✅ 消息处理成功:{message_id}")else:# 处理失败,判断是否重试ifretry_count<max_retries:# 重试:重新入队,增加重试计数headers=properties.headersor{}headers['x-retry-count']=retry_count+1channel.basic_publish(exchange='',routing_key=self.queue_name,body=body,properties=pika.BasicProperties(headers=headers,content_type=properties.content_type))channel.basic_ack(delivery_tag=method.delivery_tag)print(f"🔄 消息重试{retry_count+1}/{max_retries}:{message_id}")else:# 超过重试次数,拒绝消息(进入死信队列)channel.basic_reject(delivery_tag=method.delivery_tag,requeue=False)print(f"❌ 消息重试次数耗尽,进入死信队列:{message_id}")exceptExceptionase:# 处理异常,拒绝消息并重新入队print(f"⚠️ 消息处理异常:{e}")ifretry_count<max_retries:channel.basic_nack(delivery_tag=method.delivery_tag,requeue=True)else:channel.basic_reject(delivery_tag=method.delivery_tag,requeue=False)# 设置QoS,每次只取1条消息(公平分发)self.channel.basic_qos(prefetch_count=1)self.channel.basic_consume(queue=self.queue_name,on_message_callback=on_message)try:print(f"开始消费队列:{self.queue_name}")self.channel.start_consuming()exceptKeyboardInterrupt:print("消费停止")exceptExceptionase:print(f"消费异常:{e}")self.reconnect()self.consume_with_retry(process_func,max_retries)# 使用示例consumer=ReliableRabbitConsumer('localhost','order_queue')defprocess_order(data):"""处理订单消息"""print(f"处理订单:{data}")# 这里写业务逻辑returnTrue# 返回True表示处理成功ifconsumer.connect():consumer.consume_with_retry(process_order,max_retries=3)第二步:Kafka可靠消费(offset手动提交)
fromkafkaimportKafkaConsumer,TopicPartitionfromkafka.errorsimportCommitFailedErrorimportjsonclassReliableKafkaConsumer:"""Kafka可靠消费者:手动提交offset,保证至少一次语义"""def__init__(self,bootstrap_servers,topic,group_id):self.consumer=KafkaConsumer(topic,bootstrap_servers=bootstrap_servers,group_id=group_id,enable_auto_commit=False,# 关闭自动提交auto_offset_reset='earliest',# 从最早开始消费value_deserializer=lambdav:json.loads(v.decode('utf-8')),max_poll_records=100,# 每次最多取100条session_timeout_ms=30000,heartbeat_interval_ms=3000)print(f"✅ Kafka消费者初始化成功:{topic}")defconsume_batch_with_checkpoint(self,process_func,batch_size=100):""" 批量消费 + 手动提交offset(检查点机制) 保证:处理成功才提交,失败不提交(会重新消费) """whileTrue:# 拉取消息records=self.consumer.poll(timeout_ms=1000)ifnotrecords:continuefortp,messagesinrecords.items():batch=[]formsginmessages:batch.append(msg.value)iflen(batch)==0:continuetry:# 批量处理业务print(f"处理批次:{len(batch)}条消息")success=process_func(batch)ifsuccess:# 处理成功,手动提交offsetself.consumer.commit()print(f"✅ 批次提交成功: offset={messages[-1].offset+1}")else:# 处理失败,不提交(下批会重新消费)print(f"⚠️ 批次处理失败,等待重试")exceptExceptionase:print(f"⚠️ 批次处理异常:{e}")# 不提交offset,等待下次重新消费# 如果批次大小达到阈值,也提交iflen(batch)>=batch_size:try:self.consumer.commit()print(f"✅ 批次大小达到阈值,提交offset")exceptCommitFailedErrorase:print(f"⚠️ 提交失败(可能rebalance):{e}")# 使用示例consumer=ReliableKafkaConsumer(bootstrap_servers=['localhost:9092'],topic='order_events',group_id='rpa-order-processor')defprocess_order_batch(batch):"""批量处理订单"""fororderinbatch:print(f" 处理订单:{order['order_id']}")returnTrueconsumer.consume_batch_with_checkpoint(process_order_batch,batch_size=100)第三步:消息队列监控告警
importrequestsimporttimedefmonitor_rabbitmq_queue(host,port,username,password,queue_name,alert_threshold=1000):""" 监控RabbitMQ队列长度 队列堆积超过阈值发送告警 """# RabbitMQ Management APIurl=f"http://{host}:{port}/api/queues/%2F{queue_name}"auth=(username,password)whileTrue:try:resp=requests.get(url,auth=auth,timeout=5)ifresp.status_code==200:data=resp.json()messages=data.get('messages',0)messages_ready=data.get('messages_ready',0)messages_unack=data.get('messages_unacknowledged',0)print(f"队列{queue_name}: 总消息={messages}, 待消费={messages_ready}, 处理中={messages_unack}")ifmessages>alert_threshold:send_alert(title=f"🚨 RabbitMQ队列堆积告警",content=f"队列{queue_name}堆积{messages}条消息,超过阈值{alert_threshold}",level="high")ifmessages_unack>alert_threshold/2:send_alert(title=f"⚠️ RabbitMQ消费卡住告警",content=f"队列{queue_name}有{messages_unack}条消息处理中(未ACK),消费者可能卡住",level="medium")else:print(f"⚠️ 获取队列信息失败:{resp.status_code}")exceptExceptionase:print(f"⚠️ 监控异常:{e}")time.sleep(30)# 每30秒检查一次defmonitor_kafka_lag(bootstrap_servers,group_id,alert_threshold=1000):""" 监控Kafka消费者lag(落后消息数) lag过大说明消费速度跟不上生产速度 """fromkafka.adminimportKafkaAdminClientfromkafka.structsimportTopicPartition admin=KafkaAdminClient(bootstrap_servers=bootstrap_servers)whileTrue:try:# 用kafka-consumer-groups.sh工具查询lagimportsubprocess cmd=f"kafka-consumer-groups --bootstrap-server{bootstrap_servers[0]}--describe --group{group_id}"result=subprocess.run(cmd,shell=True,capture_output=True,text=True)ifresult.returncode==0:lines=result.stdout.strip().split('\n')forlineinlines[1:]:# 跳过表头parts=line.split()iflen(parts)>=5:topic=parts[1]partition=parts[2]lag=int(parts[4])ifparts[4]!='None'else0iflag>alert_threshold:send_alert(title=f"🚨 Kafka消费Lag告警",content=f"Group={group_id}, Topic={topic}, Partition={partition}, Lag={lag}",level="high")else:print(f"⚠️ 查询Kafka lag失败:{result.stderr}")exceptExceptionase:print(f"⚠️ 监控异常:{e}")time.sleep(30)defsend_alert(title,content,level="medium"):"""发送告警到企微/钉钉"""webhook_url="https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=YOUR_KEY"emoji="🚨"iflevel=="high"else"⚠️"full_content=f"{emoji}**{title}**\n\n{content}\n\n时间:{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}"payload={"msgtype":"markdown","markdown":{"content":full_content}}try:resp=requests.post(webhook_url,json=payload,timeout=5)ifresp.json().get('errcode')==0:print(f"✅ 告警已发送:{title}")else:print(f"⚠️ 告警发送失败:{resp.text}")exceptExceptionase:print(f"⚠️ 告警发送异常:{e}")第四步:影刀RPA完整流程编排
【启动】流程需要监听消息队列时启动 ↓ 【Python节点】consumer.connect() → 连接RabbitMQ/Kafka ↓ 【Python节点】consumer.consume_with_retry() → 开始消费 ↓ 【循环】收到消息 ↓ 【业务处理】RPA流程处理具体业务(如自动下单、数据同步) ↓ 【条件判断】处理成功? ├─ 是 → 【Python节点】channel.basic_ack() → 确认消息 └─ 否 → 【Python节点】channel.basic_nack() → 拒绝并重新入队 ↓ 【Python节点】monitor_rabbitmq_queue() → 后台监控队列堆积(线程) ↓ 【条件判断】队列堆积 > 1000? ├─ 是 → 【企微告警】发送队列堆积告警 └─ 否 → 继续消费 ↓ 【异常捕获】连接断开? ├─ 是 → 【Python节点】reconnect() → 自动重连 └─ 否 → 继续有什么坑
坑1:RabbitMQ消息丢失的经典场景
- 生产者没开confirm机制 → 消息没到broker就认为发送成功了
- 队列没设置durable=True → broker重启队列丢失
- 消费者没用手动ACK → 消息投递给消费者但还没处理broker就认为已消费
解决方案(生产者confirm机制):
# 生产者开启confirm机制channel.confirm_delivery()try:channel.basic_publish(exchange='',routing_key='my_queue',body='message',properties=pika.BasicProperties(delivery_mode=2,# 消息持久化message_id='unique-id-123'# 去重用))print("✅ 消息已确认送达broker")exceptpika.exceptions.UnroutableError:print("⚠️ 消息无法路由,需要重发或记录")坑2:Kafka重复消费
消费者处理了消息但还没提交offset就挂了,重启后会重新消费同一条消息。
解决方案:幂等性处理(业务层去重):
defprocess_with_idempotency(message):"""幂等性处理:同一条消息不会重复生效"""message_id=message.get('message_id')# 用Redis记录已处理的message_idimportredis r=redis.Redis(host='localhost',port=6379,db=0)# SET NX:只有key不存在时才设置成功(原子操作)ifr.setnx(f"msg:{message_id}","1"):r.expire(f"msg:{message_id}",86400)# 24小时过期# 第一次处理,执行业务逻辑do_business(message)returnTrueelse:# 重复消息,直接跳过print(f"⚠️ 重复消息,跳过:{message_id}")returnTrue坑3:队列堆积百万条,如何快速消费?
TEMU店群矩阵自动化运营核价报活动
单消费者处理速度跟不上,队列越堆越多。
解决方案:批量消费 + 水平扩展
# 1. 增加消费者实例(同一group_id启动多个进程/容器)# 2. 批量拉取消息(减少网络开销)# 3. 多线程处理(注意线程安全)fromconcurrent.futuresimportThreadPoolExecutor executor=ThreadPoolExecutor(max_workers=10)defconsume_concurrent(consumer,topic,group_id):"""多线程并发消费"""whileTrue:records=consumer.poll(timeout_ms=1000)fortp,messagesinrecords.items():futures=[]formsginmessages:future=executor.submit(process_message,msg.value)futures.append(future)# 等待所有线程处理完,再提交offsetforfutureinfutures:future.result()consumer.commit()坑4:测试环境和生产环境共用队列
开发测试时连错队列,把生产消息消费掉了。
解决方案:环境隔离
# 队列命名规范:{env}.{service}.{event}# 生产:prod.order.create# 测试:test.order.create# 开发:dev.order.createQUEUE_NAME=f"{ENV}.order.create"# ENV从环境变量读取总结
| 保证等级 | 手段 | 适用场景 |
|---|---|---|
| 最多一次(可能丢) | 自动ACK,不重试 | 日志收集,丢几条没关系 |
| 最少一次(可能重复) | 手动ACK + 重试 | 绝大多数业务场景(配合幂等) |
| 恰好一次(不丢不重) | 事务消息 / 幂等 + 去重表 | 支付、账务等核心场景 |
核心经验:
消费者一定要手动ACK,不要让broker自动ACK
业务逻辑必须幂等,应对重复消费
消息处理失败后不要一直重试,进死信队列人工处理
一定要监控队列堆积Lag,早发现问题早处理
编程学习
技术分享
实战经验