如何在Python中实现分布式异步并行进程的同步?
如何在Python中实现分布式异步并行进程的同步?
看起来你遇到的是典型的分布式环境下双向同步的问题——客户端工作站得等中央服务器完成电流设置才能读电压,服务器又得等客户端读完电压才能继续,对吧?结合你的代码场景,我给你几个Python里实际可行的解决方案,都是项目里常用的:
方案一:基于消息队列的双向确认(最推荐,扩展性强)
这种方式用消息队列做中间件,通过“通知-确认”的模式实现两边的同步,解耦性好,就算后续工作站数量增加也能轻松扩展。比如用RabbitMQ配合pika库实现:
服务器端代码
import pika # 初始化RabbitMQ连接(替换成你的服务器地址) connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明两个专用队列:一个发电流设置完成的通知,一个收电压读取完成的通知 channel.queue_declare(queue='current_set_done') channel.queue_declare(queue='voltage_read_done') def set_current(): # 这里写你实际的设置电流逻辑 print("服务器:完成电流设置") def wait_for_voltage_read(): # 阻塞等待客户端的读取完成通知 _, _, body = channel.basic_get(queue='voltage_read_done', auto_ack=True) if body: print("服务器:收到客户端电压读取完成的通知") # 主流程 z = 0 set_current() # 给所有客户端广播电流设置完成的消息 channel.basic_publish(exchange='', routing_key='current_set_done', body='ready_to_read') wait_for_voltage_read() y = 0 # 替换成你的实际业务值 z += y connection.close()
客户端代码
import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() channel.queue_declare(queue='current_set_done') channel.queue_declare(queue='voltage_read_done') def wait_for_current_set(): # 阻塞等待服务器的电流设置完成通知 _, _, body = channel.basic_get(queue='current_set_done', auto_ack=True) if body: print("客户端:收到服务器电流设置完成的通知") def read_voltage(): # 这里写你实际的读取电压逻辑 print("客户端:完成电压读取") return 5 # 模拟读取到的电压值 # 主流程 z = 0 wait_for_current_set() y = read_voltage() z += y # 给服务器发送电压读取完成的通知 channel.basic_publish(exchange='', routing_key='voltage_read_done', body='read_finished') connection.close()
如果觉得RabbitMQ太重,也可以用Redis的消息队列来替代,原理是一样的。
方案二:基于Redis分布式信号量(轻量简单场景适用)
如果你的场景比较简单,不需要复杂的消息路由,用Redis的分布式信号量就能搞定,依赖redis-py库:
服务器端代码
import redis # 连接Redis(替换成你的Redis地址) r = redis.Redis(host='localhost', port=6379, db=0, socket_timeout=5) def set_current(): print("服务器:完成电流设置") # 主流程 z = 0 set_current() # 释放"电流已设置"的信号,设置30秒过期避免客户端崩溃导致卡死 r.setex('current_set_flag', 30, '1') # 等待"电压已读取"的信号 while not r.exists('voltage_read_flag'): pass # 清理信号,方便下次使用 r.delete('voltage_read_flag') y = 0 z += y
客户端代码
import redis r = redis.Redis(host='localhost', port=6379, db=0, socket_timeout=5) def read_voltage(): print("客户端:完成电压读取") return 3 # 主流程 z = 0 # 等待"电流已设置"的信号 while not r.exists('current_set_flag'): pass # 清理信号 r.delete('current_set_flag') y = read_voltage() z += y # 释放"电压已读取"的信号,设置30秒过期 r.setex('voltage_read_flag', 30, '1')
这个方案好处是轻量,但要记得给信号加过期时间,防止某一方崩溃导致另一方一直阻塞。
方案三:基于RPC同步调用(标准库就能实现)
如果你的流程是严格的“服务器做A→客户端做B→服务器继续”,且是一对一的场景,用Python标准库的XML-RPC就能实现,不用额外装第三方服务:
服务器端代码
from xmlrpc.server import SimpleXMLRPCServer class SyncService: def __init__(self): self.current_is_set = False self.voltage_is_read = False self.voltage_value = 0 def finish_set_current(self): # 实际设置电流的逻辑写在这里 print("服务器:完成电流设置") self.current_is_set = True return True def wait_voltage_read(self): # 阻塞等待客户端读取完成 while not self.voltage_is_read: pass self.voltage_is_read = False return self.voltage_value def notify_voltage_read(self, value): self.voltage_value = value self.voltage_is_read = True return True # 启动RPC服务 server = SimpleXMLRPCServer(('localhost', 8000)) server.register_instance(SyncService()) print("RPC服务器启动,等待客户端连接...") server.serve_forever()
客户端代码
import xmlrpc.client # 连接服务器的RPC服务 proxy = xmlrpc.client.ServerProxy('http://localhost:8000/') def wait_current_set(): # 等待服务器完成电流设置 while not proxy.finish_set_current(): pass print("客户端:服务器已完成电流设置") def read_voltage(): value = 4 # 模拟读取的电压值 print("客户端:完成电压读取,值为", value) return value # 主流程 z = 0 wait_current_set() y = read_voltage() z += y # 通知服务器已完成电压读取 proxy.notify_voltage_read(y)
这个方案的优点是零额外依赖,用Python自带的库就能实现,但扩展性稍差,适合简单的一对一同步场景。
不管用哪种方案,都要记得加上异常处理和超时逻辑,比如网络断连、进程崩溃的情况,避免整个同步流程卡住。
备注:内容来源于stack exchange,提问作者david
相关产品推荐
相关产品推荐

