Python multiprocessing connection recv_bytes阻塞无返回问题排查
TCP流量转发程序中recv_bytes()阻塞的解决方案
问题背景
尝试用Python的multiprocessing模块实现简单TCP流量转发程序:监听指定端口,接收入站连接后与目标服务器建立出站TCP连接,随后在两个连接间双向转发原始数据。但测试时发现,调用recv_bytes()后会持续阻塞,即使连接中存在数据。
简化测试代码
from multiprocessing import Process from multiprocessing.connection import Listener, Client, wait def start_serving(listen_port, outbound_port): with Listener(('', listen_port)) as server: print(f"Waiting for connections on port {listen_port}") with server.accept() as inbound_conn: print(f"Connection accepted from {server.last_accepted}") outbound_conn = Client(('localhost', outbound_port)) print(f"Connected to port {outbound_port}") readers = [inbound_conn, outbound_conn] print(f"inbound_reader = {inbound_conn}") print(f"outbound_reader = {outbound_conn}") while readers: for r in wait(readers): try: print(f"Calling recv_bytes with reader {r}") data = r.recv_bytes() # This blocks even when there's data print(f"Out of recv_bytes with reader {r}") except EOFError: readers.remove(r) else: fwd_to_conn = None if r is inbound_conn: fwd_to_conn = outbound_conn print("read from inbound connection") elif r is outbound_conn: fwd_to_conn = outbound_conn # 此处为笔误,应指向inbound_conn print("read from outbound connection") if fwd_to_conn is not None: print(f"Forwarding {len(bytes)} bytes") # 此处为笔误,应改为len(data) fwd_to_conn.send_bytes(data) forwarder = Process(target=start_serving, daemon=True, args=(19001, 19002)) forwarder.start() forwarder.join()
测试步骤
在两个终端分别执行:
$ echo "Hi from outbound" | nc -l -p 19002
$ echo "Hi from inbound" | nc localhost 19001
运行输出
Waiting for connections on port 19001 Connection accepted from ('127.0.0.1', 56874) Connected to port 19002 inbound_reader = <multiprocessing.connection.Connection object at 0x7ff3730e2850> outbound_reader = <multiprocessing.connection.Connection object at 0x7ff3730e2a60> Calling recv_bytes with reader <multiprocessing.connection.Connection object at 0x7ff3730e2850>
问题原因
查看multiprocessing.connection.Connection源码后发现,recv_bytes()方法要求消息的前4字节为消息长度标识,而我们需要转发的是无格式的原始TCP数据,没有这4字节的长度头,导致recv_bytes()一直等待符合格式的数据,从而阻塞。
解决方案
multiprocessing.connection.Connection本质是基于socket封装的,我们可以直接操作其底层socket对象,调用原生recv()方法读取原始字节流,无需解析格式。修改后的代码如下:
from multiprocessing import Process from multiprocessing.connection import Listener, Client, wait def start_serving(listen_port, outbound_port): with Listener(('', listen_port)) as server: print(f"Waiting for connections on port {listen_port}") with server.accept() as inbound_conn: print(f"Connection accepted from {server.last_accepted}") outbound_conn = Client(('localhost', outbound_port)) print(f"Connected to port {outbound_port}") # 获取底层socket对象 inbound_sock = inbound_conn._handle outbound_sock = outbound_conn._handle readers = [inbound_sock, outbound_sock] # 建立socket到目标转发socket的映射 conn_map = { inbound_sock: outbound_sock, outbound_sock: inbound_sock } print(f"inbound_sock = {inbound_sock}") print(f"outbound_sock = {outbound_sock}") while readers: for r in wait(readers): try: print(f"Calling recv with socket {r}") # 读取原始数据,每次最多读取4096字节 data = r.recv(4096) if not data: # 收到空字节表示连接已关闭 readers.remove(r) continue print(f"Received {len(data)} bytes from socket {r}") except Exception as e: print(f"Error reading from socket: {e}") readers.remove(r) else: fwd_sock = conn_map[r] try: # 使用sendall确保所有数据发送完成 fwd_sock.sendall(data) print(f"Forwarded {len(data)} bytes to socket {fwd_sock}") except Exception as e: print(f"Error forwarding data: {e}") # 一侧连接出错,关闭两侧连接 readers.remove(fwd_sock) readers.remove(r) # 手动关闭连接 inbound_conn.close() outbound_conn.close() forwarder = Process(target=start_serving, daemon=True, args=(19001, 19002)) forwarder.start() forwarder.join()
关键修改点
- 直接调用
Connection对象的_handle属性获取底层socket,用recv(4096)读取原始字节流,规避recv_bytes()的格式限制 - 建立socket映射关系,简化双向转发逻辑
- 使用
sendall()确保数据完整发送 - 完善连接异常和关闭的处理逻辑
同时修正了原代码中的两处笔误:
- 从outbound_conn读取数据时,转发目标改为
inbound_conn - 打印转发字节数时,将
len(bytes)修正为len(data)
内容的提问来源于stack exchange,提问作者Eugene S
相关产品推荐
相关产品推荐

