影刀RPA 生产者消费者模式:采集和处理分离的流程设计
影刀RPA 生产者消费者模式采集和处理分离的流程设计你有没有遇到过这种情况采集很快处理很慢网页上2秒采下来的数据写入ERP要30秒。如果串行执行——采集→处理→采集→处理——大部分时间都在等处理完成。生产者消费者模式把采集和处理拆开采集的只采集处理的只处理中间用队列缓冲。模式原理生产者采集 缓冲区 消费者处理 ↓ ↓ ↓ 采数据 ──放入──→ [队列] ──取出──→ 写ERP 采数据 ──放入──→ [队列] ──取出──→ 写ERP 采数据 ──放入──→ [队列] ──取出──→ 写ERP ...等着... 采完 ──告诉──→ ↓ ──收到──→ 结束生产者和消费者各跑各的生产者不用等消费者处理完才能采下一条。在影刀中的三种实现方式方式一文件队列最简单用文件夹作为缓冲区。生产者把采集结果写成文件消费者轮询文件夹读取处理。店群矩阵自动化突破运营极限生产者流程采集生产者 collect.flow 打开网页 → 翻页循环采集 FOR EACH 批次: 采集当前页数据 # 存成队列文件 文件名 fqueue_{batch_id:05d}.json 【写入JSON文件】→ queue文件夹  【日志输出】f已生产批次 {batch_id} batch_id 1 END FOR # 生产完毕放一个哨兵文件 【写入文件】→ queue_done.txt内容done消费者流程数据处理消费者 process.flow processed 0 WHILE True: # 获取队列文件夹中的所有数据文件 【获取文件列表】→ queue文件夹 → *.json → file_list IF file_list 为空: # 检查是否生产已完成 IF queue_done.txt 存在: 【日志输出】所有数据已处理完毕 BREAK ELSE: 【等待 3 秒】 # 等生产者产出新数据 CONTINUE # 取第一个文件处理 【读取JSON文件】→ file_list[0] 处理数据写ERP等... processed 记录数 # 处理完删除或移动到已处理文件夹 【移动文件】→ 移动到 processed/ 文件夹 END WHILE方式二数据库队列用数据库表作为队列生产者INSERT消费者SELECTDELETE。# 生产者cursor.execute( INSERT INTO task_queue (data, status, created_at) VALUES (%s, pending, NOW()) ,(json.dumps(batch_data),))conn.commit()# 消费者cursor.execute( SELECT id, data FROM task_queue WHERE status pending ORDER BY id ASC LIMIT 1 FOR UPDATE -- 锁定这一行防止其他消费者也读到 )rowcursor.fetchone()ifrow:datajson.loads(row[1])# 处理data...cursor.execute(UPDATE task_queue SET status done WHERE id %s,(row[0],))conn.commit()FOR UPDATE是行级锁——如果有多个消费者同时读队列只有一个能锁住这一行其他消费者跳过这一行去读下一行。避免了同一条数据被重复处理。方式三主流程多开子流程伪并行影刀不支持真正的多线程并行但可以做到伪并行——主流程同时触发两个子流程但它们还是串行执行。真正的生产者消费者在影刀里的最佳实践是拆成两个独立的流程用文件或数据库做桥梁。流程A生产者每天9点启动 → 采集数据 → 写入queue文件 流程B消费者每天9点05启动 → 轮询queue文件 → 处理数据两个流程之间的5分钟间隔确保数据已经生产完毕。队列文件的命名与锁定文件队列最大的风险生产者在写入文件到一半时消费者就开始读了。解决方案1先写临时文件写完了再重命名。importosimportjsonimporttempfile queue_dirrC:\RPA_data\queuetemp_dirrC:\RPA_data\temp# 先写到临时文件temp_pathos.path.join(temp_dir,ftemp_{batch_id}.json)withopen(temp_path,w,encodingutf-8)asf:json.dump(data,f,ensure_asciiFalse)# 写完后原子性移动同一磁盘上的rename是原子操作final_pathos.path.join(queue_dir,fqueue_{batch_id:05d}.json)os.replace(temp_path,final_path)# 原子性替换os.replace()在同一磁盘分区上是原子操作——消费者要么看到空文件还没move要么看到完整文件move完成不会看到写了一半的文件。解决方案2用处理标记文件。生产者写完后额外创建一个.ready标记文件temu店群自动化报活动案例生产者 写入 queue_001.json 写入 queue_001.ready ← 消费者只处理有.ready标记的文件 消费者 只处理同时存在 .json 和 .ready 两个文件的条目消费者处理失败的恢复消费者处理失败时不能让这条数据丢失# 消费者处理失败重试max_retries3forfile_pathinqueue_files:retries0whileretriesmax_retries:try:dataread_queue_file(file_path)process_data(data)mark_as_done(file_path)break# 处理成功exceptExceptionase:retries1print(f处理失败{retries}/{max_retries}{e})ifretriesmax_retries:move_to_failed(file_path)# 移到失败文件夹else:# 全部重试都失败print(f文件{file_path}处理失败已移至failed文件夹)监控队列积压生产者和消费者之间的速度不匹配——如果消费者太慢队列文件越积越多。需要监控# 在消费者里加积压监控importos queue_dirrC:\RPA_data\queuepending_countlen([fforfinos.listdir(queue_dir)iff.endswith(.json)])ifpending_count100:print(f⚠️ 队列积压严重{pending_count}个文件待处理)# 可以发告警通知elifpending_count50:print(f队列积压{pending_count}个文件)总结生产者消费者模式解决的是生产速度和消费速度不匹配的问题。文件队列最简单实用数据库队列支持多消费者并发。核心要点先写临时文件再重命名防止读到半成品处理失败要有重试和死信队列要监控队列积压防止数据堆积。作者林焱