如何并发写入多个TCP套接字?二进制转JSON流场景实践需求
看起来你已经找对方向了——用多进程处理独立的TCP流确实是提升并发导入效率的好方案。我来帮你把代码补全优化,确保每个流都能稳定高效地运行:
完整并发TCP流发送示例
首先给你一个可以直接运行的完整代码,包含完善的异常处理和资源管理:
import socket import sys import multiprocessing as mp import json def write_tcp_stream(host, port, packet, max_retries=2): sock = None for attempt in range(max_retries + 1): try: # 创建并连接TCP套接字 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.connect((host, port)) # 将解析后的二进制数据序列化为JSON并编码为字节流 json_bytes = json.dumps(packet).encode('utf-8') # 确保数据完整发送(sendall会自动处理分片) sock.sendall(json_bytes) # 如果需要确认服务器接收,可以添加响应接收逻辑 # response = sock.recv(1024) # print(f"服务器响应: {response.decode('utf-8')}") return except socket.error as msg: sys.stderr.write(f"[ERROR] 第{attempt+1}次发送失败: {str(msg)}\n") if sock: sock.close() finally: if sock: sock.close() sys.stderr.write(f"[ERROR] 超过最大重试次数,放弃发送该数据包\n") def main(): # 服务器配置 server_host = "127.0.0.1" server_port = 9000 # 这里替换成你解析好的二进制数据包列表 parsed_packets = [ {"id": 1, "binary_data": "parsed_content_1"}, {"id": 2, "binary_data": "parsed_content_2"}, {"id": 3, "binary_data": "parsed_content_3"} ] # 设置并发进程数:IO密集型任务可以适当高于CPU核心数 worker_count = mp.cpu_count() * 2 # 使用进程池管理并发任务,自动处理进程创建/销毁 with mp.Pool(worker_count) as pool: # 用starmap传递多参数任务 pool.starmap( write_tcp_stream, [(server_host, server_port, pkt) for pkt in parsed_packets] ) if __name__ == "__main__": main()
关键细节说明
- 资源安全管理:用
try-finally确保套接字无论成功失败都会关闭,避免系统资源泄漏 - 重试机制:添加了简单的重试逻辑,应对临时网络波动
- 进程池选型:
multiprocessing.Pool比手动创建Process更省心,自动平衡任务分配 - 并发数调整:IO密集型场景(比如TCP发送)可以把进程数设为CPU核心数的2-4倍,充分利用网络带宽
额外优化建议
- 边解析边发送:如果解析二进制数据也耗时,可以把解析逻辑放到
write_tcp_stream函数里,避免主进程一次性解析所有数据占用大量内存 - 日志替代stderr:用Python的
logging模块替代sys.stderr.write,方便后续排查问题 - 批量发送控制:如果数据包数量极大,可以分批次发送,避免进程池瞬间创建过多进程压垮服务器
内容的提问来源于stack exchange,提问作者Dang Khoa
相关产品推荐
相关产品推荐

