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

使用DuckDB处理Parquet数据时内存占用过高及OOM问题咨询

问题:DuckDB分组聚合内存占用过高及优化疑问

场景与任务

我正在使用DuckDB处理存储于Hive式目录结构(按年、月、日、小时分区)的Parquet文件,每个文件包含约150列且需全量读取。当前任务需基于flowKey、condition2、condition3三列对相似行进行分组聚合,将匹配的多行合并为一行,涉及约147列的聚合操作,待分组的总组数为3990717。

遇到的问题

在配备32GB内存、8核的机器上,仅处理1GB的Parquet数据就消耗了约20GB内存,且出现Out of Memory Error: Failed to allocate block of Bytes错误,想了解该内存消耗水平在DuckDB的此类操作中是否正常。

已尝试的优化措施

  • 执行conn.execute("SET temp_directory = '/tmp/duckdb_swap';")开启磁盘溢出,但未起到作用
  • 尝试先分离无需分组的唯一行,再用PyArrow将其与分组结果合并

聚合查询示例

其余列采用类似聚合方式:

select 
    first(sourceIPAddress) as sourceIPAddress,    
    first(destinationIPAddress) as destinationIPAddress, 
    sum(goodputBytesIn) as goodputBytesIn, 
    sum(goodputBytesOut) as goodputBytesOut, 
    avg(duration) as duration, 
    cast(sum(cwrCount) as int64) as cwrCount, 
    cast(sum(rstCount) as int64) as rstCount, 
    cast(sum(synCount) as int64) as synCount, 
    cast(sum(finCount) as int64) as finCount, 
    cast(sum(dst2srcFinCount) as int64) as dst2srcFinCount, 
    avg(byteDistMean) as byteDistMean, 
    avg(byteDistStd) as byteDistStd, 
    avg(entropy) as entropy,
    avg(totalEntropy) as totalEntropy, 
    max(pktlenStoCmax) as pktlenStoCmax, 
    avg(pktlenStoCstd) as pktlenStoCstd, 
    FIRST(app) FILTER (WHERE app <> 'Unknown') AS app,
    FIRST(httpCode) FILTER (WHERE httpCode <> -1) AS httpCode,
    FIRST(userAgent) FILTER (WHERE userAgent IS NOT NULL) AS userAgent,
    FIRST(hostname) FILTER (WHERE hostname IS NOT NULL) AS hostname
from 
    arrow_table 
group by 
    flowKey, condition2, condition3

数据处理流程

从Parquet文件加载数据,生成condition2和condition3两列并追加至PyArrow表,以此作为查询输入,利用DuckDB与PyArrow的高效零拷贝集成特性。

DuckDB连接管理类

class duckdbManager:
    """
    IMPLEMENTING CONTEXT MANAGER FOR HANDLING CONNECTION'S IN DUCKDB
    RETURNS: CONNECTION OBJECT FOR DUCKDB
    """
    def __init__(self):
        pass

    def __enter__(self):
        self.connection = duckdb.connect(database=':memory:', read_only=False)
        return self.connection

    def __exit__(self, exc_type, exc_value, traceback):
        if self.connection:
            self.connection.close()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 23:03:21