影刀RPA 消息队列自动化:RabbitMQ Kafka可靠性保证
影刀RPA 消息队列自动化RabbitMQ Kafka可靠性保证什么情况用什么 → 怎么做 → 有什么坑作者林焱 | 飞行社出品什么情况用什么用RPA处理业务时需要和生产系统的消息队列对接——要么从队列取任务要么推送任务。但消息丢了、重复消费了排查起来要命。这套方案适合RPA流程和消息队列集成异步任务处理保证消息不丢、不重复处理监控队列堆积情况及时处理核心工具影刀RPA pikaRabbitMQ kafka-pythonKafka 监控告警怎么做拼多多店群自动化报活动上架第一步RabbitMQ可靠消费手动ACK 重试importpikaimportjsonfromdatetimeimportdatetimeimporttimeclassReliableRabbitConsumer:RabbitMQ可靠消费者保证消息不丢失def__init__(self,host,queue_name,usernameguest,passwordguest):self.hosthost self.queue_namequeue_name self.credentialspika.PlainCredentials(username,password)self.connectionNoneself.channelNonedefconnect(self):建立连接带重连机制try:parameterspika.ConnectionParameters(hostself.host,credentialsself.credentials,heartbeat600,# 心跳超时10分钟blocked_connection_timeout300)self.connectionpika.BlockingConnection(parameters)self.channelself.connection.channel()# 声明队列幂等已存在不会报错self.channel.queue_declare(queueself.queue_name,durableTrue,# 队列持久化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_retries3): 可靠消费手动ACK 失败重试 process_func: 业务处理函数返回True表示处理成功 defon_message(channel,method,properties,body):message_idproperties.message_idormethod.delivery_tag retry_countproperties.headers.get(x-retry-count,0)ifproperties.headerselse0try:# 解析消息ifproperties.content_typeapplication/json:datajson.loads(body)else:databody.decode(utf-8)print(f收到消息{message_id}:{str(data)[:100]})# 处理业务successprocess_func(data)ifsuccess:# 处理成功手动ACKchannel.basic_ack(delivery_tagmethod.delivery_tag)print(f✅ 消息处理成功:{message_id})else:# 处理失败判断是否重试ifretry_countmax_retries:# 重试重新入队增加重试计数headersproperties.headersor{}headers[x-retry-count]retry_count1channel.basic_publish(exchange,routing_keyself.queue_name,bodybody,propertiespika.BasicProperties(headersheaders,content_typeproperties.content_type))channel.basic_ack(delivery_tagmethod.delivery_tag)print(f 消息重试{retry_count1}/{max_retries}:{message_id})else:# 超过重试次数拒绝消息进入死信队列channel.basic_reject(delivery_tagmethod.delivery_tag,requeueFalse)print(f❌ 消息重试次数耗尽进入死信队列:{message_id})exceptExceptionase:# 处理异常拒绝消息并重新入队print(f⚠️ 消息处理异常:{e})ifretry_countmax_retries:channel.basic_nack(delivery_tagmethod.delivery_tag,requeueTrue)else:channel.basic_reject(delivery_tagmethod.delivery_tag,requeueFalse)# 设置QoS每次只取1条消息公平分发self.channel.basic_qos(prefetch_count1)self.channel.basic_consume(queueself.queue_name,on_message_callbackon_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)# 使用示例consumerReliableRabbitConsumer(localhost,order_queue)defprocess_order(data):处理订单消息print(f处理订单:{data})# 这里写业务逻辑returnTrue# 返回True表示处理成功ifconsumer.connect():consumer.consume_with_retry(process_order,max_retries3)第二步Kafka可靠消费offset手动提交fromkafkaimportKafkaConsumer,TopicPartitionfromkafka.errorsimportCommitFailedErrorimportjsonclassReliableKafkaConsumer:Kafka可靠消费者手动提交offset保证至少一次语义def__init__(self,bootstrap_servers,topic,group_id):self.consumerKafkaConsumer(topic,bootstrap_serversbootstrap_servers,group_idgroup_id,enable_auto_commitFalse,# 关闭自动提交auto_offset_resetearliest,# 从最早开始消费value_deserializerlambdav:json.loads(v.decode(utf-8)),max_poll_records100,# 每次最多取100条session_timeout_ms30000,heartbeat_interval_ms3000)print(f✅ Kafka消费者初始化成功:{topic})defconsume_batch_with_checkpoint(self,process_func,batch_size100): 批量消费 手动提交offset检查点机制 保证处理成功才提交失败不提交会重新消费 whileTrue:# 拉取消息recordsself.consumer.poll(timeout_ms1000)ifnotrecords:continuefortp,messagesinrecords.items():batch[]formsginmessages:batch.append(msg.value)iflen(batch)0:continuetry:# 批量处理业务print(f处理批次:{len(batch)}条消息)successprocess_func(batch)ifsuccess:# 处理成功手动提交offsetself.consumer.commit()print(f✅ 批次提交成功: offset{messages[-1].offset1})else:# 处理失败不提交下批会重新消费print(f⚠️ 批次处理失败等待重试)exceptExceptionase:print(f⚠️ 批次处理异常:{e})# 不提交offset等待下次重新消费# 如果批次大小达到阈值也提交iflen(batch)batch_size:try:self.consumer.commit()print(f✅ 批次大小达到阈值提交offset)exceptCommitFailedErrorase:print(f⚠️ 提交失败可能rebalance:{e})# 使用示例consumerReliableKafkaConsumer(bootstrap_servers[localhost:9092],topicorder_events,group_idrpa-order-processor)defprocess_order_batch(batch):批量处理订单fororderinbatch:print(f 处理订单:{order[order_id]})returnTrueconsumer.consume_batch_with_checkpoint(process_order_batch,batch_size100)第三步消息队列监控告警importrequestsimporttimedefmonitor_rabbitmq_queue(host,port,username,password,queue_name,alert_threshold1000): 监控RabbitMQ队列长度 队列堆积超过阈值发送告警 # RabbitMQ Management APIurlfhttp://{host}:{port}/api/queues/%2F{queue_name}auth(username,password)whileTrue:try:resprequests.get(url,authauth,timeout5)ifresp.status_code200:dataresp.json()messagesdata.get(messages,0)messages_readydata.get(messages_ready,0)messages_unackdata.get(messages_unacknowledged,0)print(f队列{queue_name}: 总消息{messages}, 待消费{messages_ready}, 处理中{messages_unack})ifmessagesalert_threshold:send_alert(titlef RabbitMQ队列堆积告警,contentf队列{queue_name}堆积{messages}条消息超过阈值{alert_threshold},levelhigh)ifmessages_unackalert_threshold/2:send_alert(titlef⚠️ RabbitMQ消费卡住告警,contentf队列{queue_name}有{messages_unack}条消息处理中未ACK消费者可能卡住,levelmedium)else:print(f⚠️ 获取队列信息失败:{resp.status_code})exceptExceptionase:print(f⚠️ 监控异常:{e})time.sleep(30)# 每30秒检查一次defmonitor_kafka_lag(bootstrap_servers,group_id,alert_threshold1000): 监控Kafka消费者lag落后消息数 lag过大说明消费速度跟不上生产速度 fromkafka.adminimportKafkaAdminClientfromkafka.structsimportTopicPartition adminKafkaAdminClient(bootstrap_serversbootstrap_servers)whileTrue:try:# 用kafka-consumer-groups.sh工具查询lagimportsubprocess cmdfkafka-consumer-groups --bootstrap-server{bootstrap_servers[0]}--describe --group{group_id}resultsubprocess.run(cmd,shellTrue,capture_outputTrue,textTrue)ifresult.returncode0:linesresult.stdout.strip().split(\n)forlineinlines[1:]:# 跳过表头partsline.split()iflen(parts)5:topicparts[1]partitionparts[2]lagint(parts[4])ifparts[4]!Noneelse0iflagalert_threshold:send_alert(titlef Kafka消费Lag告警,contentfGroup{group_id}, Topic{topic}, Partition{partition}, Lag{lag},levelhigh)else:print(f⚠️ 查询Kafka lag失败:{result.stderr})exceptExceptionase:print(f⚠️ 监控异常:{e})time.sleep(30)defsend_alert(title,content,levelmedium):发送告警到企微/钉钉webhook_urlhttps://qyapi.weixin.qq.com/cgi-bin/webhook/send?keyYOUR_KEYemojiiflevelhighelse⚠️full_contentf{emoji}**{title}**\n\n{content}\n\n时间:{datetime.now().strftime(%Y-%m-%d %H:%M:%S)}payload{msgtype:markdown,markdown:{content:full_content}}try:resprequests.post(webhook_url,jsonpayload,timeout5)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() → 自动重连 └─ 否 → 继续有什么坑坑1RabbitMQ消息丢失的经典场景生产者没开confirm机制 → 消息没到broker就认为发送成功了队列没设置durableTrue → broker重启队列丢失消费者没用手动ACK → 消息投递给消费者但还没处理broker就认为已消费解决方案生产者confirm机制# 生产者开启confirm机制channel.confirm_delivery()try:channel.basic_publish(exchange,routing_keymy_queue,bodymessage,propertiespika.BasicProperties(delivery_mode2,# 消息持久化message_idunique-id-123# 去重用))print(✅ 消息已确认送达broker)exceptpika.exceptions.UnroutableError:print(⚠️ 消息无法路由需要重发或记录)坑2Kafka重复消费消费者处理了消息但还没提交offset就挂了重启后会重新消费同一条消息。解决方案幂等性处理业务层去重defprocess_with_idempotency(message):幂等性处理同一条消息不会重复生效message_idmessage.get(message_id)# 用Redis记录已处理的message_idimportredis rredis.Redis(hostlocalhost,port6379,db0)# SET NX只有key不存在时才设置成功原子操作ifr.setnx(fmsg:{message_id},1):r.expire(fmsg:{message_id},86400)# 24小时过期# 第一次处理执行业务逻辑do_business(message)returnTrueelse:# 重复消息直接跳过print(f⚠️ 重复消息跳过:{message_id})returnTrue坑3队列堆积百万条如何快速消费TEMU店群矩阵自动化运营核价报活动单消费者处理速度跟不上队列越堆越多。解决方案批量消费 水平扩展# 1. 增加消费者实例同一group_id启动多个进程/容器# 2. 批量拉取消息减少网络开销# 3. 多线程处理注意线程安全fromconcurrent.futuresimportThreadPoolExecutor executorThreadPoolExecutor(max_workers10)defconsume_concurrent(consumer,topic,group_id):多线程并发消费whileTrue:recordsconsumer.poll(timeout_ms1000)fortp,messagesinrecords.items():futures[]formsginmessages:futureexecutor.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_NAMEf{ENV}.order.create# ENV从环境变量读取总结保证等级手段适用场景最多一次可能丢自动ACK不重试日志收集丢几条没关系最少一次可能重复手动ACK 重试绝大多数业务场景配合幂等恰好一次不丢不重事务消息 / 幂等 去重表支付、账务等核心场景核心经验消费者一定要手动ACK不要让broker自动ACK业务逻辑必须幂等应对重复消费消息处理失败后不要一直重试进死信队列人工处理一定要监控队列堆积Lag早发现问题早处理