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

如何并发写入多个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倍,充分利用网络带宽
额外优化建议
  1. 边解析边发送:如果解析二进制数据也耗时,可以把解析逻辑放到write_tcp_stream函数里,避免主进程一次性解析所有数据占用大量内存
  2. 日志替代stderr:用Python的logging模块替代sys.stderr.write,方便后续排查问题
  3. 批量发送控制:如果数据包数量极大,可以分批次发送,避免进程池瞬间创建过多进程压垮服务器

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:11:53