You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 15:43:14