使用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操作。当然,如果解压缩耗时过长,也可能间接导致套接字闲置,但从你的情况看,核心问题还是处理速度跟不上下载速度。
解决方案
优先选择先下载后处理
直接把文件下载到本地,再读取本地文件进行处理,这是最稳妥的方案,完全避免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)若必须流式处理
- 调整超时设置:将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()
- 调整超时设置:将requests的
优化MariaDB写入性能
- 使用批量插入:攒够N行数据后执行一次INSERT语句(比如每1000行一次),减少DB交互次数。
- 关闭自动提交:手动控制事务提交,比如每处理10000行提交一次事务。
- 使用
LOAD DATA INFILE:将处理好的数据写入本地临时文件,再用该命令批量导入,速度比逐行插入快几个数量级。
内容的提问来源于stack exchange,提问作者audiodude
相关产品推荐
相关产品推荐

