Python中实现字节流缓冲FIFO队列,优化异步数据处理与上传
解决方案:Python实现并行流式读写(带固定FIFO缓冲)
核心思路
用Python内置的queue.Queue实现固定大小的FIFO缓冲,配合多线程让数据读取、格式处理、GCS写入三个环节并行执行:
- 读线程:从数据库流式拉取数据,缓冲满时自动阻塞,避免内存占用过高
- 写线程:从缓冲取数据,缓冲空时自动阻塞,保证数据不丢失
- 格式处理:通过自定义函数灵活适配CSV/JSON或其他格式,完全脱离特定工具包限制
代码实现(通用框架)
import queue import threading from google.cloud import storage # 固定缓冲大小,按需调整(比如1000条数据) BUFFER_CAPACITY = 1000 data_buffer = queue.Queue(maxsize=BUFFER_CAPACITY) stop_signal = threading.Event() def stream_from_db(): """数据库流式读取线程""" # 替换为你的实际数据库流式读取逻辑(比如cursor.fetchmany循环) while not stop_signal.is_set(): batch = fetch_db_batch() if not batch: break for row in batch: # 缓冲满时自动阻塞,直到有空闲位置 data_buffer.put(row, block=True) def process_and_upload(): """数据处理+GCS写入线程""" client = storage.Client() bucket = client.get_bucket("your-target-bucket") blob = bucket.blob("output.csv") # 可替换为output.json # 初始化文件(比如写CSV表头) with blob.open("w") as f: write_file_header(f) # 循环处理直到读取完成且缓冲为空 while not stop_signal.is_set() or not data_buffer.empty(): try: # 缓冲空时阻塞,超时后检查停止信号 row = data_buffer.get(block=True, timeout=1) # 自定义格式转换:CSV行/JSON字符串都可 processed = format_row(row) f.write(f"{processed}\n") data_buffer.task_done() except queue.Empty: continue # -------------------------- # 以下为自定义逻辑示例,按需替换 # -------------------------- def fetch_db_batch(): """模拟数据库批量读取,替换为实际查询""" # 示例:返回100条模拟数据 return [{"id": i, "content": f"data_{i}"} for i in range(100)] def write_file_header(file): """写入CSV表头,JSON可跳过此步骤""" file.write("id,content\n") def format_row(row): """将数据库行转为CSV格式,可替换为json.dumps(row)""" return f"{row['id']},{row['content']}" # -------------------------- # 启动并管理线程 # -------------------------- read_thread = threading.Thread(target=stream_from_db) write_thread = threading.Thread(target=process_and_upload) read_thread.start() write_thread.start() # 等待读取完成,触发停止信号后等待写入线程处理剩余数据 read_thread.join() stop_signal.set() write_thread.join()
关键特性说明
- 自动阻塞控制:
queue.Queue原生实现满缓冲写阻塞、空缓冲读阻塞,无需手动处理锁和等待逻辑 - 纯并行IO:读、写线程独立运行,彻底消除串行等待的耗时问题
- 格式完全自定义:
format_row函数可自由实现CSV拼接、JSON序列化或其他格式转换,不依赖任何第三方格式工具包 - 优雅停机:通过
stop_signal确保所有缓冲数据处理完成后再退出,避免数据丢失
进阶优化方向
- 批量写入GCS:积累固定数量数据后再批量写入(比如每500条写一次),减少GCS API调用次数
- 多处理线程:如果格式转换是CPU密集型,可改用
multiprocessing.Queue配合多进程 - 错误重试:添加数据库断连、GCS上传失败的异常捕获与重试逻辑
内容的提问来源于stack exchange,提问作者Rob Allsopp
相关产品推荐
相关产品推荐

