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

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()确保数据完整发送
  • 完善连接异常和关闭的处理逻辑

同时修正了原代码中的两处笔误:

  1. 从outbound_conn读取数据时,转发目标改为inbound_conn
  2. 打印转发字节数时,将len(bytes)修正为len(data)

内容的提问来源于stack exchange,提问作者Eugene S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 22:15:39