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

如何优化QuestDB的ILP文件上传脚本?解决管道破裂问题

解决Socket大文件上传Broken Pipe问题

问题根源

Broken pipe错误本质是连接被对方关闭,常见原因:

  • 发送速度远超接收端处理能力,导致接收端主动断开连接
  • 网络波动引发连接中断
  • 原脚本逐行发送无流量控制,也没有重连/重试机制

改进后的脚本

针对18GB大文件,需实现分块传输、错误重连、流量控制,代码如下:

import socket
import sys
import time

def create_socket():
    sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    sock.settimeout(30)  # 设置超时,避免无限等待
    return sock

def send_file_chunk(sock, chunk):
    total_sent = 0
    chunk_len = len(chunk)
    while total_sent < chunk_len:
        try:
            sent = sock.send(chunk[total_sent:])
            if sent == 0:
                raise socket.error("连接已关闭")
            total_sent += sent
        except socket.error as e:
            # 发送失败,抛出异常让上层处理重连
            raise e

if __name__ == '__main__':
    server_addr = ('localhost', 9009)
    file_path = "YOUR_FILE"
    chunk_size = 4096 * 16  # 64KB分块,可根据网络调整
    retry_limit = 5  # 最大重试次数
    current_retry = 0

    while current_retry < retry_limit:
        sock = None
        try:
            sock = create_socket()
            sock.connect(server_addr)
            print(f"成功连接到服务器 {server_addr}")

            with open(file_path, 'rb') as infile:
                # 先发送文件总大小(可选,方便接收端校验)
                file_size = infile.seek(0, 2)
                infile.seek(0)
                sock.sendall(file_size.to_bytes(8, byteorder='big'))

                while True:
                    chunk = infile.read(chunk_size)
                    if not chunk:
                        break  # 文件发送完成
                    send_file_chunk(sock, chunk)
                    # 可选:等待接收端确认(如果接收端会回复ACK)
                    # ack = sock.recv(1)
                    # if ack != b'OK':
                    #     raise Exception("接收端未确认")

            print("文件上传完成")
            break  # 成功发送,退出循环

        except socket.error as e:
            print(f"传输错误: {e},正在重试 ({current_retry+1}/{retry_limit})")
            current_retry += 1
            time.sleep(2)  # 重试前等待2秒
        finally:
            if sock:
                sock.close()

    if current_retry >= retry_limit:
        sys.stderr.write(f"已达到最大重试次数 {retry_limit},上传失败")

关键改进点

  • 分块读取:用固定大小分块替代逐行读取,适配大文件,避免内存占用过高
  • 重连机制:连接中断后自动重试,最多尝试5次
  • 完整发送保证:send_file_chunk函数确保每个分块完全发送,避免部分数据丢失
  • 超时设置:防止连接或发送时无限挂起
  • 文件大小前置发送:可选功能,方便接收端提前知道文件总大小,校验完整性

额外建议

  • 如果接收端支持,添加确认机制:每发送一个分块后等待接收端的ACK信号,彻底避免发送过快导致的断开
  • 根据网络带宽调整chunk_size:带宽高可增大分块,带宽低则减小
  • 记录上传进度:可以添加进度条或日志,跟踪大文件上传状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 17:42:42