Python多进程Pipe异常排查:通信数小时后进程挂起问题
进程间Pipe通信挂起问题分析与解决
问题描述
有两个进程:核心进程运行主循环,WebSocket进程通过API获取最新价格,二者通过Pipe通信。核心进程发送请求后等待WebSocket进程的响应,运行数小时后核心进程发送请求后挂起,一直等待响应。检查发现WebSocket进程仍在运行无异常,但未接收到核心进程的请求。复现代码如下:
from time import sleep from multiprocessing import Pipe, Process class websocket_stuff: def __init__(self, pipe: Pipe): self.process_pipe = pipe # Websocket functions are not added since its not necessary for demonstration def WS_run_forever(self): # here i do the WS infinite loop, and as part of the loop i do the following: while True: if self.process_pipe.poll() is True: rec = self.process_pipe.recv() if (rec is not None) and (isinstance(rec, dict) is True): if ('request' in rec) and (rec['request'] == 'candle'): self.process_pipe.send(self.process_pipe, {'message': 'candle', 'ticker': rec['ticker'], 'data': 'OHLC data'}) class core: def __init__(self): self.name = 'ticker' self.WS_pipe, kid = Pipe() self.WS_obj = websocket_stuff(kid) self.WS_process = Process(target=self.WS_obj.WS_run_forever) self.WS_process.start() def get_data(self): self.WS_pipe.send({'request': 'candle', 'ticker': self.name}) twe = 0 while True: # It hangs/freezes here if self.WS_pipe.poll(2) is True: rec = self.WS_pipe.recv() break else: twe += 1 if twe >= 3: self.WS_pipe.send({'request': 'candle', 'ticker': self.name}) sleep(0.0001) if ('message' in rec) and (rec['message'] == 'candle') and (rec['ticker'] == self.name): return rec['data'] return None if __name__ == '__main__': c = core() while True: data = c.get_data() print(data) sleep(1)
问题原因
- Pipe.send()参数错误:WebSocket进程中调用
self.process_pipe.send(self.process_pipe, {...})时多传入了管道对象本身。Pipe的send()方法仅需传入要发送的数据对象,多余的参数会导致发送失败(实际运行中可能被异常捕获静默处理),使得核心进程的请求无法被正常响应,进而引发后续的请求堆积。 - 管道缓冲区溢出导致阻塞:核心进程在超时后重复发送请求,而WebSocket进程未处理这些请求,导致管道的发送缓冲区被填满。在默认阻塞模式下,核心进程再次调用
send()时会被挂起,表现为程序卡住。 - 缺乏异常处理机制:WebSocket进程的循环中没有异常捕获逻辑,若某次接收/发送操作抛出异常,可能导致进程逻辑卡壳,无法继续处理后续请求,但进程本身仍处于运行状态。
解决方案
修复send()参数错误:修改WebSocket进程中的发送代码,移除多余的管道参数:
# 原错误代码 # self.process_pipe.send(self.process_pipe, {'message': 'candle', 'ticker': rec['ticker'], 'data': 'OHLC data'}) # 修正后 self.process_pipe.send({'message': 'candle', 'ticker': rec['ticker'], 'data': 'OHLC data'})添加异常捕获与日志:在WebSocket进程的循环中加入异常处理,及时发现并处理错误,避免静默失败:
def WS_run_forever(self): while True: try: if self.process_pipe.poll() is True: rec = self.process_pipe.recv() if rec is not None and isinstance(rec, dict): if 'request' in rec and rec['request'] == 'candle': self.process_pipe.send({'message': 'candle', 'ticker': rec['ticker'], 'data': 'OHLC data'}) except Exception as e: # 记录异常日志,便于排查问题 print(f"WebSocket进程通信异常: {str(e)}") sleep(0.001) # 适当增加休眠,减少CPU占用优化核心进程请求逻辑:避免无限制重复发送请求,超时后检查WebSocket进程状态,必要时重启进程:
def get_data(self): max_retries = 3 retry_count = 0 while retry_count < max_retries: self.WS_pipe.send({'request': 'candle', 'ticker': self.name}) if self.WS_pipe.poll(2): rec = self.WS_pipe.recv() if 'message' in rec and rec['message'] == 'candle' and rec['ticker'] == self.name: return rec['data'] retry_count += 1 # 重试多次失败后,重启WebSocket进程 print("WebSocket无响应,尝试重启进程") self.WS_process.terminate() self.WS_process.join() self.WS_process = Process(target=self.WS_obj.WS_run_forever) self.WS_process.start() return None降低CPU占用:在两个进程的循环中添加适当的休眠时间(如
sleep(0.001)),避免空循环导致CPU使用率过高,间接影响通信稳定性。
内容的提问来源于stack exchange,提问作者user12358312
相关产品推荐
相关产品推荐

