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

UNIX Socket通信问题:解决消息长度未知导致客户端阻塞

解决UNIX Socket未知消息长度的接收阻塞问题

你的问题核心是TCP(包括UNIX Stream Socket)是流式协议,没有内置的消息边界,客户端硬编码固定接收长度(16字节),当服务器返回的结果字节数不足时,recv()会一直阻塞等待更多数据,导致客户端陷入死循环。

下面提供三种实用的解决思路,以及对应修改后的代码:

方法1:使用消息边界标记(推荐,适合文本消息)

在服务器发送的结果末尾添加一个特殊标记(比如换行符\n),客户端接收数据时,直到读取到这个标记就停止,无需预设长度。

修改后的server.py

import numpy as np
import socket
import os

from signal import signal, SIGPIPE, SIG_DFL
signal(SIGPIPE, SIG_DFL)

def main():
    server_address = "/tmp/uds_socket"

    # Remove old sockets
    try:
        os.unlink(server_address)
    except OSError:
        if os.path.exists(server_address):
            raise

    server_socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    print(f"Loading the model: {server_address}")
    server_socket.bind(server_address)
    server_socket.listen(1)

    while True:
        print("Waiting for a connection...")
        client_socket, client_address = server_socket.accept()
        print("Connection established...")

        try:
            while True:
                msg = client_socket.recv(128)
                if msg:
                    data = list(map(float, msg.decode().split(',')))
                    print(f"Received: {data}, with length: {len(msg)}")
                    result = sum(data)
                    # 发送结果时添加换行符作为边界
                    client_socket.sendall(f"{result}\n".encode())
                else:
                    print("No data")
                    break
        finally:
            print("Closing the connection")
            client_socket.close()

if __name__ == "__main__":
    main()

修改后的client.py

import numpy as np
import socket
import errno
import time
import sys
import os

def main():
    client_socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    server_address = "/tmp/uds_socket"

    print(f"Connecting to: {server_address}")
    try:
        client_socket.connect(server_address)
    except:
        print("Error")
        sys.exit(1)

    while True:
        data = np.random.uniform(low=-10.0, high=10.0, size=4).tolist()
        msg = ','.join(map(str, data))

        try:
            client_socket.sendall(msg.encode())
            
            # 接收直到遇到换行符
            result_buffer = b''
            while True:
                chunk = client_socket.recv(16)  # 每次接收16字节提升效率
                if not chunk:
                    break
                if b'\n' in chunk:
                    # 拆分出换行符前的内容并结束循环
                    result_buffer += chunk.split(b'\n')[0]
                    break
                result_buffer += chunk
            print(f"Received sum from server: {result_buffer.decode()}")

        except IOError as e:
            if e.errno == errno.EPIPE:
                print(f"Error here: {e.errno}")
        except KeyboardInterrupt:
            break

        time.sleep(1)

    client_socket.close()

if __name__ == "__main__":
    main()

方法2:先发送消息长度前缀(更可靠,适合二进制/复杂消息)

服务器先发送结果的字节长度(用固定4字节的大端序整数),客户端先读取长度,再根据长度读取对应字节数的内容。

修改后的server.py

import numpy as np
import socket
import os
import struct

from signal import signal, SIGPIPE, SIG_DFL
signal(SIGPIPE, SIG_DFL)

def main():
    server_address = "/tmp/uds_socket"

    try:
        os.unlink(server_address)
    except OSError:
        if os.path.exists(server_address):
            raise

    server_socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    print(f"Loading the model: {server_address}")
    server_socket.bind(server_address)
    server_socket.listen(1)

    while True:
        print("Waiting for a connection...")
        client_socket, client_address = server_socket.accept()
        print("Connection established...")

        try:
            while True:
                msg = client_socket.recv(128)
                if msg:
                    data = list(map(float, msg.decode().split(',')))
                    print(f"Received: {data}, with length: {len(msg)}")
                    result_str = str(sum(data))
                    result_bytes = result_str.encode()
                    # 先发送4字节的长度前缀(大端序)
                    length_prefix = struct.pack('!I', len(result_bytes))
                    client_socket.sendall(length_prefix + result_bytes)
                else:
                    print("No data")
                    break
        finally:
            print("Closing the connection")
            client_socket.close()

if __name__ == "__main__":
    main()

修改后的client.py

import numpy as np
import socket
import errno
import time
import sys
import os
import struct

def main():
    client_socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    server_address = "/tmp/uds_socket"

    print(f"Connecting to: {server_address}")
    try:
        client_socket.connect(server_address)
    except:
        print("Error")
        sys.exit(1)

    while True:
        data = np.random.uniform(low=-10.0, high=10.0, size=4).tolist()
        msg = ','.join(map(str, data))

        try:
            client_socket.sendall(msg.encode())
            
            # 先接收4字节的长度前缀
            length_data = b''
            while len(length_data) < 4:
                chunk = client_socket.recv(4 - len(length_data))
                if not chunk:
                    break
                length_data += chunk
            if len(length_data) != 4:
                print("Invalid length prefix")
                continue
            result_length = struct.unpack('!I', length_data)[0]
            
            # 再接收对应长度的结果
            result_buffer = b''
            while len(result_buffer) < result_length:
                chunk = client_socket.recv(result_length - len(result_buffer))
                if not chunk:
                    break
                result_buffer += chunk
            print(f"Received sum from server: {result_buffer.decode()}")

        except IOError as e:
            if e.errno == errno.EPIPE:
                print(f"Error here: {e.errno}")
        except KeyboardInterrupt:
            break

        time.sleep(1)

    client_socket.close()

if __name__ == "__main__":
    main()

方法3:单次请求后关闭连接(简单但性能一般)

客户端发送请求后,服务器处理完就关闭连接,客户端接收直到recv()返回空字节,就知道消息结束。这种方式不需要额外的边界处理,但会频繁创建/销毁连接,适合低并发场景。

修改后的server.py(核心修改)

在发送结果后主动跳出循环,触发连接关闭:

# ... 原有代码 ...
result = sum(data)
client_socket.sendall(str(result).encode())
# 发送完成后跳出循环,触发finally块关闭连接
break
# ... 原有代码 ...

修改后的client.py(核心修改)

# ... 原有代码 ...
try:
    client_socket.sendall(msg.encode())
    
    result_buffer = b''
    while True:
        chunk = client_socket.recv(16)
        if not chunk:
            break
        result_buffer += chunk
    print(f"Received sum from server: {result_buffer.decode()}")
    # 关闭当前连接,下次循环重新建立连接
    client_socket.close()
    client_socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    client_socket.connect(server_address)
# ... 原有代码 ...

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:54:53