如何在Python单TCP Socket上实现同时收发消息及接收Callbacks?
基于TCP的心跳客户端优雅实现方案
针对你的需求,推荐使用Python标准库的selectors模块实现多路复用IO,既能避免冗余的超时处理,又能轻松实现消息接收的回调机制。以下是具体实现思路和代码示例:
核心实现思路
- 多路复用IO:用
selectors.DefaultSelector监听TCP Socket的可读事件,同时利用select的超时参数触发定时心跳发送,无需单独维护复杂的超时逻辑。 - 回调机制:将消息接收后的业务逻辑抽离为独立的回调函数,在读取到服务器消息时直接调用,实现IO操作与业务逻辑的解耦。
- 非阻塞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
相关产品推荐
相关产品推荐

