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

FastAPI服务处理Trino数据并上传S3后内存无法释放的问题求助

FastAPI服务处理Trino数据并上传S3后内存无法释放的问题求助

大家好,我这边碰到一个棘手的内存问题,想请教下社区的朋友们。我们运行着一个FastAPI服务,流程是从Trino拉取数据,用PyArrow处理后上传到AWS S3的Parquet数据集里,但每次请求处理完后内存都降不下来,哪怕做了各种手动清理操作都不管用。

架构概述

  • 框架:FastAPI
  • 数据源:Trino
  • 处理工具:PyArrow、Polars
  • 存储:AWS S3(Parquet格式)

处理流程

  1. API接收用户请求
  2. 从Trino拉取目标数据
  3. 将数据转换为PyArrow Table格式
  4. 通过write_to_dataset上传到S3并按字段分区
  5. 尝试手动释放占用的内存

核心代码片段

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 10:28:07