Python多进程管道实现句子分词任务技术求助
多进程管道通信(匿名管道+FIFO)代码优化方案
问题背景
需要实现三个进程间的通信:
- 源进程:从元组读取消息,写入匿名管道
- 转换进程:读取匿名管道内容,分词后将每个单词单独一行写入FIFO
- 输出进程:读取FIFO内容并打印
以下是针对现有代码的优化建议及修正后的实现:
核心问题修正与优化点
1. 修复管道逻辑错误
原代码错误创建了冗余的transformer_pipe,且转换进程误将写入该管道作为逻辑目标,实际应读取源进程的匿名管道。需删除冗余管道,让转换进程直接对接源进程的输出管道。
2. 修复消息读取阻塞问题
源进程写入的消息未添加换行符,导致转换进程的readline()会持续阻塞(直到管道关闭),无法逐条处理消息。需在每条消息末尾添加换行符,确保readline()能正确识别单条消息边界。
3. 优化文件描述符管理
- 使用
with语句管理文件对象,自动完成资源关闭,替代手动try-finally的close()操作,代码更简洁且避免资源泄漏。 - 确保每个进程仅保留所需的文件描述符,减少不必要的资源占用。
4. 完善FIFO创建逻辑
创建FIFO前先清理残留的同名管道,避免因历史残留导致启动报错;同时针对FIFO操作捕获特定异常,提升代码鲁棒性。
5. 细化异常处理
针对FIFO创建、管道读写等操作捕获特定OSError,而非笼统捕获所有Exception,便于精准定位问题。
修正后的完整代码
import os import sys import time MESSAGES = ( b"I love Python programming", b"Inter-process communication is important", b"Pipes and FIFOs are used for communication", ) def source_process(write_pipe): with write_pipe: for message in MESSAGES: # 添加换行符,确保转换进程能通过readline正确读取单条消息 write_pipe.write(message + b'\n') write_pipe.flush() time.sleep(1) # 延迟便于观察输出 def transformer_process(read_pipe, fifo_path): with read_pipe: # 打开FIFO写端,会阻塞直到读端被打开 try: fifo_fd = os.open(fifo_path, os.O_WRONLY) except OSError as e: print(f"Failed to open FIFO for writing: {e}", file=sys.stderr) return try: while True: message = read_pipe.readline() if not message: break # 管道关闭,退出循环 # 分词并转换为每行一个单词的格式 words = message.strip().split() if words: transformed = b'\n'.join(words) + b'\n' try: os.write(fifo_fd, transformed) except OSError as e: if e.errno == 32: # 管道破裂,输出进程已退出 break finally: os.close(fifo_fd) def output_process(fifo_path): # 打开FIFO读端,会阻塞直到写端被打开 try: fifo_fd = os.open(fifo_path, os.O_RDONLY) except OSError as e: print(f"Failed to open FIFO for reading: {e}", file=sys.stderr) return try: while True: data = os.read(fifo_fd, 1024) if not data: break # FIFO写端关闭,退出循环 print(data.decode(), end='', flush=True) finally: os.close(fifo_fd) def main(): fifo_path = "my_fifo" # 先清理残留的FIFO try: os.unlink(fifo_path) except FileNotFoundError: pass try: os.mkfifo(fifo_path) except OSError as e: print(f"Failed to create FIFO: {e}", file=sys.stderr) return # 创建源进程与转换进程间的匿名管道 source_pipe_read, source_pipe_write = os.pipe() # 启动源进程 source_pid = os.fork() if source_pid == 0: # 子进程关闭不需要的读端 os.close(source_pipe_read) source_process(os.fdopen(source_pipe_write, 'wb')) sys.exit(0) else: # 父进程关闭不需要的写端 os.close(source_pipe_write) # 启动转换进程 transformer_pid = os.fork() if transformer_pid == 0: # 子进程关闭不需要的写端(父进程已关闭,冗余但安全) os.close(source_pipe_write) transformer_process(os.fdopen(source_pipe_read, 'rb'), fifo_path) sys.exit(0) else: # 父进程关闭不需要的读端 os.close(source_pipe_read) # 启动输出进程 output_pid = os.fork() if output_pid == 0: output_process(fifo_path) sys.exit(0) # 等待所有子进程完成 for pid in [source_pid, transformer_pid, output_pid]: os.waitpid(pid, 0) # 清理FIFO try: os.unlink(fifo_path) except FileNotFoundError: pass if __name__ == "__main__": main()
内容的提问来源于stack exchange,提问作者chris kadas
相关产品推荐
相关产品推荐

