FastAPI服务处理Trino数据并上传S3后内存无法释放的问题求助
FastAPI服务处理Trino数据并上传S3后内存无法释放的问题求助
大家好,我这边碰到一个棘手的内存问题,想请教下社区的朋友们。我们运行着一个FastAPI服务,流程是从Trino拉取数据,用PyArrow处理后上传到AWS S3的Parquet数据集里,但每次请求处理完后内存都降不下来,哪怕做了各种手动清理操作都不管用。
架构概述
- 框架:FastAPI
- 数据源:Trino
- 处理工具:PyArrow、Polars
- 存储:AWS S3(Parquet格式)
处理流程
- API接收用户请求
- 从Trino拉取目标数据
- 将数据转换为PyArrow Table格式
- 通过
write_to_dataset上传到S3并按字段分区 - 尝试手动释放占用的内存
核心代码片段
import pyarrow as pa import pyarrow.parquet as pq import pyarrow.fs as pafs # 从Trino拉取并处理数据 data = fetch_data_from_trino() arrow_table = pa.table(list(zip(*data))) # 初始化S3文件系统,配置为高性能写入 s3_filesystem = pafs.S3FileSystem() # 上传到S3,按organisation字段分区 pq.write_to_dataset( arrow_table, root_path=f"{RESULT_STORAGE_BUCKET}/{s3_storage_path}", partition_cols=['organisation'], filesystem=s3_filesystem ) # 尝试手动释放内存 del data del arrow_table
已尝试的解决方法及观察结果
1. 手动删除大对象
我们在上传完成后用del删掉了原始数据列表data和PyArrow Table对象arrow_table,但内存占用依然居高不下。
某次请求的内存统计数据:
"resource_stats": { "memory_mb": { "start": 140.3004568, "peak": 589.02921865, "end": 587.01258 } }
2. 强制触发垃圾回收
在del操作之后,我们调用了gc.collect(),甚至写了工具函数多次触发各代GC,还用psutil监控内存变化,但内存几乎没有下降:
import os import psutil import gc from typing import Dict def force_garbage_collection_and_cleanup() -> Dict[str, int]: try: # 记录清理前的内存占用 memory_before = psutil.Process(os.getpid()).memory_info().rss / (1024 ** 2) # 多次触发垃圾回收,确保清理彻底 objects_collected = 0 for _ in range(3): collected = gc.collect() objects_collected += collected # 记录清理后的内存占用 memory_after = psutil.Process(os.getpid()).memory_info().rss / (1024 ** 2) return { "memory_before_mb": round(memory_before, 2), "memory_after_mb": round(memory_after, 2), "objects_collected": objects_collected } except Exception as cleanup_error: logger.warning(f"垃圾回收失败: {str(cleanup_error)}") return {}
执行后的内存统计:
"resource_stats": { "memory_mb": { "start": 190.56, "peak": 579.6, "end": 579.7 } }
3. 自定义深度内存释放函数
我们还写了一个更激进的内存释放函数:不仅分代多次触发GC,还在Linux服务器上调用malloc_trim(0)尝试让glibc把闲置内存还给系统,但最终内存释放效果依然不明显。
import gc import psutil import platform import ctypes from typing import Dict def force_memory_release() -> Dict: """ 强制Python释放内存给操作系统: 1. 多次触发垃圾回收清理无引用对象 2. Linux环境下调用malloc_trim()让glibc归还内存 3. 统计并返回内存变化情况 """ # 记录清理前内存 process = psutil.Process() memory_before_mb = process.memory_info().rss / 1024 / 1024 # 分代触发垃圾回收 collected_gen0 = gc.collect(0) # 年轻代 collected_gen1 = gc.collect(1) # 中年代 collected_gen2 = gc.collect(2) # 老年代 collected_final = gc.collect() total_objects_collected = collected_gen0 + collected_gen1 + collected_gen2 + collected_final # Linux环境下调用malloc_trim归还内存 malloc_trim_success = False if platform.system() == 'Linux': try: libc = ctypes.CDLL("libc.so.6") malloc_trim = libc.malloc_trim malloc_trim.argtypes = [ctypes.c_size_t] malloc_trim.restype = ctypes.c_int result = malloc_trim(0) # 0表示归还所有可释放的内存 malloc_trim_success = bool(result) print(f"malloc_trim调用{'成功' if malloc_trim_success else '失败,无内存可归还'}") except Exception as e: print(f"调用malloc_trim出错: {e}") malloc_trim_success = False # 记录清理后内存 memory_after_mb = process.memory_info().rss / 1024 / 1024 memory_freed_mb = max(0, memory_before_mb - memory_after_mb) cleanup_stats = { 'objects_collected': total_objects_collected, 'memory_before_mb': round(memory_before_mb, 2), 'memory_after_mb': round(memory_after_mb, 2), 'memory_freed_mb': round(memory_freed_mb, 2), 'malloc_trim_success': malloc_trim_success, 'platform': platform.system() } print( f"内存清理完成: 回收对象数{total_objects_collected}, " f"释放内存{memory_freed_mb:.2f}MB ({memory_before_mb:.2f}MB -> {memory_after_mb:.2f}MB), " f"malloc_trim状态: {'成功' if malloc_trim_success else '不可用/失败'}" ) return cleanup_stats # 调用释放函数 force_memory_release()
想请教的问题
有没有朋友遇到过类似的场景?尤其是PyArrow和FastAPI结合时的内存泄漏问题?我们怀疑是不是PyArrow的某些内部对象(比如S3FileSystem实例、或者write_to_dataset的临时对象)没有被正确回收?或者有没有其他我们没考虑到的内存优化点?
内容来源于stack exchange
相关产品推荐
相关产品推荐

