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

如何在Python单TCP Socket上实现同时收发消息及接收Callbacks?

基于TCP的心跳客户端优雅实现方案

针对你的需求,推荐使用Python标准库的selectors模块实现多路复用IO,既能避免冗余的超时处理,又能轻松实现消息接收的回调机制。以下是具体实现思路和代码示例:

核心实现思路

  1. 多路复用IO:用selectors.DefaultSelector监听TCP Socket的可读事件,同时利用select的超时参数触发定时心跳发送,无需单独维护复杂的超时逻辑。
  2. 回调机制:将消息接收后的业务逻辑抽离为独立的回调函数,在读取到服务器消息时直接调用,实现IO操作与业务逻辑的解耦。
  3. 非阻塞Socket:将Socket设置为非阻塞模式,避免单线程下IO操作阻塞整个程序运行。

完整代码示例

import selectors
import socket
import time

class TCPClient:
    def __init__(self, host, port, heartbeat_interval=0.7, message_callback=None):
        self.host = host
        self.port = port
        self.heartbeat_interval = heartbeat_interval
        self.last_heartbeat_time = 0
        self.message_callback = message_callback or self.default_message_callback
        
        # 初始化selector和socket
        self.selector = selectors.DefaultSelector()
        self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        self.sock.setblocking(False)
        self.sock.connect_ex((host, port))
        
        # 注册可读事件
        self.selector.register(self.sock, selectors.EVENT_READ, self._handle_read)

    def default_message_callback(self, message):
        """默认消息回调,可自定义替换"""
        print(f"收到服务器消息: {message.decode('utf-8').strip()}")

    def _handle_read(self, sock, mask):
        """处理socket可读事件,读取消息并调用回调"""
        try:
            data = sock.recv(1024)
            if data:
                self.message_callback(data)
            else:
                print("服务器关闭连接")
                self.selector.unregister(sock)
                sock.close()
        except Exception as e:
            print(f"读取消息出错: {e}")
            self.selector.unregister(sock)
            sock.close()

    def _send_heartbeat(self):
        """发送心跳消息"""
        try:
            self.sock.sendall(b"heartbeat")
            self.last_heartbeat_time = time.time()
            print("已发送心跳")
        except Exception as e:
            print(f"发送心跳出错: {e}")
            self.selector.unregister(self.sock)
            self.sock.close()

    def run(self):
        """客户端主循环"""
        while True:
            # 计算距离下次心跳的剩余时间
            now = time.time()
            time_to_heartbeat = self.heartbeat_interval - (now - self.last_heartbeat_time)
            timeout = max(0, time_to_heartbeat)
            
            # 监听事件,超时后触发心跳检查
            events = self.selector.select(timeout=timeout)
            for key, mask in events:
                callback = key.data
                callback(key.fileobj, mask)
            
            # 检查是否需要发送心跳
            if time.time() - self.last_heartbeat_time >= self.heartbeat_interval:
                self._send_heartbeat()

if __name__ == "__main__":
    # 自定义消息回调示例
    def custom_message_handler(msg):
        print(f"自定义处理消息: {msg.decode('utf-8').strip()}")
    
    client = TCPClient("127.0.0.1", 8080, message_callback=custom_message_handler)
    try:
        client.run()
    except KeyboardInterrupt:
        print("客户端退出")
        client.selector.close()
        client.sock.close()

代码关键点说明

  • 多路复用处理:通过selectors.select()同时监听Socket的可读事件和心跳超时,无需单独设置Socket的超时时间,避免了冗余的超时判断代码。
  • 回调机制:初始化客户端时可传入自定义的message_callback函数,当读取到服务器消息时自动调用该函数,完全解耦了IO读取和业务逻辑。
  • 非阻塞模式:Socket设置为非阻塞后,connect_ex()、sendall()等操作不会阻塞主循环,确保心跳和消息接收能同时进行。
  • 心跳控制:通过记录上次心跳时间,结合select的超时参数,精准控制每0.7秒发送一次心跳,同时避免频繁发送。

优势对比

相比直接使用socket模块手动处理超时,该方案的优势在于:

  • 代码结构更清晰,无需在多个地方重复处理超时逻辑
  • 天然支持并发IO操作,单线程即可同时处理消息接收和心跳发送
  • 回调机制让业务逻辑更易维护和扩展

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 18:03:18