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

使用Python requests流式处理大BZ2文件时遇连接中断问题求助

流式读取大文件时出现连接中断错误

我有一个3.3GB的大型数据文件,尝试使用requests的stream=True和iter_content()流式读取其内容并逐行处理。逻辑本身能正常运行,但处理15-20分钟后会出现以下错误:

<SNIP>
urllib3.exceptions.IncompleteRead: IncompleteRead(123551405 bytes read, 3479106984 more expected)

The above exception was the direct cause of the following exception:

<SNIP>
urllib3.exceptions.ProtocolError: ('Connection broken: IncompleteRead(123551405 bytes read, 3479106984 more expected)', IncompleteRead(123551405 bytes read, 3479106984 more expected))

During handling of the above exception, another exception occurred:

<SNIP>
requests.exceptions.ChunkedEncodingError: ('Connection broken: IncompleteRead(123551405 bytes read, 3479106984 more expected)', IncompleteRead(123551405 bytes read, 3479106984 more expected))

我的代码如下:

@contextmanager
def get_pageview_response():
  url = get_pageview_url()
  with requests.get(url,
                    stream=True,
                    headers={
                        'User-Agent': WP1_USER_AGENT,
                        'Connection': 'keep-alive'
                    },
                    timeout=120) as r:
    
    r.raise_for_status()
    yield r


def raw_pageviews(decode=False):

  def as_bytes():
    with get_pageview_response() as r:
      decompressor = BZ2Decompressor()
      trailing = b''
      # Read data in 128 MB chunks
      for http_chunk in r.iter_content(chunk_size=128 * 1024 * 1024):
        data = decompressor.decompress(http_chunk)
        lines = [line for line in data.split(b'\n') if line]
        if not lines:
          continue

        # Reunite incomplete lines
        yield trailing + lines[0]
        yield from lines[1:-1]
        trailing = lines[-1]

      # Nothing left, yield the last line
      yield trailing

  if decode:
    for line in as_bytes():
      yield line.decode('utf-8')
  else:
    yield from as_bytes()

def pageview_components():
  for line in raw_pageviews():
    pass  # Processing logic goes here

我的疑问

  • 猜测是处理数据(含写入MariaDB的I/O操作)速度跟不上HTTP流,导致HTTP套接字闲置,服务器关闭了连接。这个猜测是否正确?
  • 该问题是否与我的BZ2解压缩步骤有关?
  • 已尝试将HTTP块大小从16MB增加到128MB来扩大缓冲区,但没有效果。

解答

你的猜测基本正确:当数据处理(尤其是MariaDB写入这类I/O密集型操作)的速度远慢于HTTP下载速度时,连接会因为长时间没有数据交互(套接字闲置)被服务器或中间网络网关主动断开。

关于BZ2解压缩:它属于CPU密集型操作,但如果你的瓶颈在DB写入,那么解压缩不是直接原因——它只是完成下载后的解码步骤,真正阻塞流程的是后续的I/O操作。当然,如果解压缩耗时过长,也可能间接导致套接字闲置,但从你的情况看,核心问题还是处理速度跟不上下载速度。

解决方案

  1. 优先选择先下载后处理
    直接把文件下载到本地,再读取本地文件进行处理,这是最稳妥的方案,完全避免HTTP连接超时问题,还能随时暂停、重试,不占用服务器连接资源。示例下载代码:

    import requests
    
    def download_file(url, local_path):
        with requests.get(url, stream=True, headers={'User-Agent': WP1_USER_AGENT}) as r:
            r.raise_for_status()
            with open(local_path, 'wb') as f:
                for chunk in r.iter_content(chunk_size=128*1024*1024):
                    f.write(chunk)
    
  2. 若必须流式处理

    • 调整超时设置:将requests的timeout参数设为None或更大的值(比如300),注意部分服务器会有自身的超时限制,这个参数仅控制客户端超时。
    • 启用TCP Keep-Alive:通过配置urllib3的套接字参数,定期发送心跳包保持连接,避免服务器判定连接闲置:
      import socket
      from requests.adapters import HTTPAdapter
      from urllib3.util.retry import Retry
      import requests
      
      session = requests.Session()
      # 添加重试机制,应对临时连接中断
      retry = Retry(total=5, backoff_factor=1, allowed_methods=["GET"])
      adapter = HTTPAdapter(max_retries=retry)
      session.mount("https://", adapter)
      # 配置TCP Keep-Alive参数
      session.adapters['https://'].poolmanager.connection_pool_kw['socket_options'] = [
          (socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1),
          (socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 60),  # 60秒后开始发送心跳
          (socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 10), # 每10秒发一次心跳
          (socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 5)     # 连续5次失败则断开
      ]
      
    • 异步解耦下载与处理:用多线程分离下载、解压和处理流程,用队列缓冲数据,让下载流程不被处理速度拖累:
      import queue
      import threading
      import bz2
      
      def download_worker(url, q):
          with requests.get(url, stream=True, headers={'User-Agent': WP1_USER_AGENT}) as r:
              r.raise_for_status()
              decompressor = bz2.BZ2Decompressor()
              trailing = b''
              for chunk in r.iter_content(chunk_size=128*1024*1024):
                  data = decompressor.decompress(chunk)
                  lines = data.split(b'\n')
                  if trailing:
                      lines[0] = trailing + lines[0]
                  trailing = lines.pop() if lines else b''
                  for line in lines:
                      if line:
                          q.put(line)
              if trailing:
                  q.put(trailing)
          q.put(None)  # 发送结束信号
      
      def process_worker(q):
          while True:
              line = q.get()
              if line is None:
                  break
              # 替换为你的实际处理逻辑(如写入MariaDB)
              pass
      
      def run_stream_processing():
          q = queue.Queue(maxsize=1000)  # 限制队列大小,防止内存溢出
          url = get_pageview_url()
          download_thread = threading.Thread(target=download_worker, args=(url, q))
          process_thread = threading.Thread(target=process_worker, args=(q,))
          download_thread.start()
          process_thread.start()
          download_thread.join()
          process_thread.join()
      
  3. 优化MariaDB写入性能

    • 使用批量插入:攒够N行数据后执行一次INSERT语句(比如每1000行一次),减少DB交互次数。
    • 关闭自动提交:手动控制事务提交,比如每处理10000行提交一次事务。
    • 使用LOAD DATA INFILE:将处理好的数据写入本地临时文件,再用该命令批量导入,速度比逐行插入快几个数量级。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 14:24:53