影刀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):![在这里插入图片描述](https://i-blog.csdnimg.cn/direct/04607c8b07bc41d7a32bb93365b82a04.png#pic_center)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():![在这里插入图片描述](https://i-blog.csdnimg.cn/direct/caf08607117d478faa136a1ccd8f4846.png#pic_center)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早发现问题早处理

相关新闻

《收支日历图》四、ArkTS日历开发避坑指南

《收支日历图》四、ArkTS日历开发避坑指南

ArkTS 日历开发避坑指南:10 个高频问题与修复方案 摘要:在使用 HarmonyOS ArkUI 开发每日收支日历图的过程中,笔者踩了不少坑——从 JavaScript Date 的月份索引陷阱,到多月份数据过滤遗漏,再到状态管理 V2 的装饰器误…

2026/7/21 6:18:27 阅读更多 →
嵌入式LCD驱动时序配置详解:从TFT/STN原理到TI寄存器实战

嵌入式LCD驱动时序配置详解:从TFT/STN原理到TI寄存器实战

1. 项目概述与核心价值 在嵌入式系统开发中,驱动一块LCD屏幕远不止是“点亮”那么简单。很多工程师在拿到一块新的LCD模组时,常常会卡在初始化配置这一步,明明按照手册写了驱动代码,屏幕要么一片漆黑,要么显示错位、闪…

2026/7/21 6:18:27 阅读更多 →
使用libtcc实现C语言动态编译与JIT技术

使用libtcc实现C语言动态编译与JIT技术

1. 为什么需要将C编译器嵌入程序?在传统开发模式中,C代码需要预先编译成可执行文件才能运行。但某些场景下,我们需要更灵活的代码执行方式:动态脚本功能:让用户输入C代码片段即时执行插件系统扩展:允许第三…

2026/7/21 6:17:26 阅读更多 →

最新新闻

神经网络:通用函数逼近器

神经网络:通用函数逼近器

在《[[AI 研究方法的演变]]》那篇笔记中,我们沿着研究方法的演变脉络,理解了 AI 当前主流的研究为什么会走向深度神经网络。具体来说就是:在逻辑符号无法对所有规则进行编码,而概率方法又卡在了特征工程的情况下。深度神经网络提供…

2026/7/22 1:46:31 阅读更多 →
如何快速保护硬件隐私:EASY-HWID-SPOOFER硬件信息修改终极指南

如何快速保护硬件隐私:EASY-HWID-SPOOFER硬件信息修改终极指南

如何快速保护硬件隐私:EASY-HWID-SPOOFER硬件信息修改终极指南 【免费下载链接】EASY-HWID-SPOOFER 基于内核模式的硬件信息欺骗工具 项目地址: https://gitcode.com/gh_mirrors/ea/EASY-HWID-SPOOFER 在数字时代,你的电脑硬件信息就像数字指纹一…

2026/7/22 1:46:31 阅读更多 →
深入解析TI eHRPWM死区生成与故障保护模块的配置与调试

深入解析TI eHRPWM死区生成与故障保护模块的配置与调试

1. 项目概述:为什么我们需要关注eHRPWM的“内功”?在电力电子和电机驱动的世界里,PWM(脉冲宽度调制)就像是驱动系统的“心跳”。无论是让电机平稳旋转,还是让电源高效转换,都离不开精准的PWM信号…

2026/7/22 1:46:31 阅读更多 →
linux系统移植pjsua库实现sip通话功能 一、概述

linux系统移植pjsua库实现sip通话功能 一、概述

本文实现pjsua开源库的交叉编译以及通过调用pjsua库api实现sip语言双向通话功能,包括注册上线sip服务器、接收处理指令(电话邀请、挂断等)、双向音频对讲功能。 二、交叉编译pjsua库 1.解压 tar -zxvf 2.15.1.tar.gz; cd pjproject-2.15.1; 2.配置编译选项 ./config…

2026/7/22 1:46:31 阅读更多 →
B站视频下载神器终极指南:轻松获取4K大会员专属内容

B站视频下载神器终极指南:轻松获取4K大会员专属内容

B站视频下载神器终极指南:轻松获取4K大会员专属内容 【免费下载链接】bilibili-downloader B站视频下载,支持下载大会员清晰度4K,持续更新中 项目地址: https://gitcode.com/gh_mirrors/bil/bilibili-downloader 想要随时随地观看B站上…

2026/7/22 1:46:31 阅读更多 →
RAG技术解析:大语言模型与知识检索的融合应用

RAG技术解析:大语言模型与知识检索的融合应用

1. RAG技术核心解析:当大语言模型遇上知识检索检索增强生成(Retrieval-Augmented Generation,简称RAG)正在重塑AI内容生成的技术范式。这项技术的本质是将大语言模型(LLM)的生成能力与精准的信息检索系统相…

2026/7/22 1:45:30 阅读更多 →

日新闻

TI DSP系统配置模块SYSCFG详解:中断机制与主设备优先级配置实战

TI DSP系统配置模块SYSCFG详解:中断机制与主设备优先级配置实战

1. 项目概述与SYSCFG模块的核心价值在嵌入式系统,尤其是像TI C6000系列这样的高性能DSP开发中,我们常常会与芯片手册里那些密密麻麻的寄存器打交道。很多开发者可能更关注算法实现、内存优化或者外设驱动,但对于一个稳定、高效的系统而言&…

2026/7/22 0:00:26 阅读更多 →
微信Server酱:高到达率的应急通知方案实践

微信Server酱:高到达率的应急通知方案实践

1. 为什么我们需要"最次"的通知方案? 在数字化协作环境中,消息通知系统的重要性不言而喻明。但现实情况是,企业级通知方案往往需要复杂的API对接(如企业微信、钉钉、飞书),个人开发者的小项目又经…

2026/7/22 0:00:26 阅读更多 →
甲方要的“简洁“PPT,到底是简洁还是省事?

甲方要的“简洁“PPT,到底是简洁还是省事?

甲方说"简洁一点",乙方听到的是"少做几页"。甲方说"不要太复杂",乙方理解成"别放图表了"。结果交过去,甲方说"我说的简洁不是这个意思"。"简洁"这个词在PPT语境里,是…

2026/7/22 0:00:26 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/21 8:48:31 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/21 5:34:47 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/21 8:25:39 阅读更多 →

月新闻