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

如何按指定列值拆分DataFrame输出为Parquet并保证分区纯净?

如何用Dask实现按key列的纯净分区输出Parquet

我明白你想要的是严格按key列的值进行分区,确保每个分区只包含单一key值,而不是仅仅把相同key的数据尽量凑在一起。你之前用的shuffle(['key'])之所以达不到效果,是因为它本质是哈希shuffle:它会把相同key的数据分到同一个分区,但无法保证一个分区里只有一个key(极端情况下不同key可能哈希碰撞),而且如果单个key的数据量超过你设置的partition_size,还会被拆成多个分区。

下面给你两种可靠的实现方案,根据你的后续使用场景选择:

方案1:通过指定divisions实现纯净分区(适合Dask后续加载)

这种方法会让Dask的DataFrame每个分区严格对应一个key值,输出的Parquet是一个整体数据集,后续用Dask加载时可以快速按key筛选。

步骤如下:

  1. 先计算所有唯一的key值并排序,用来生成分区边界
  2. 设置索引时指定这些边界,确保每个分区只包含单个key

完整代码示例:

import dask.dataframe as dd

# 读取原始数据
df = dd.read_csv(
    "/Users/ecerulm/Downloads/test/**/*.txt.gz",
    include_path_column=True,
    sep="\t",
    compression='gzip',
    blocksize=None
)

# 生成key列(和你原来的逻辑一致)
df['basename'] = df.path.str.rpartition('/')[2]
df['MAC'] = df.basename.str.partition('.')[2]
df['MAC'] = df.MAC.str.partition('.')[0]
df['key'] = df.MAC.str[-2:]

# 1. 获取所有唯一的key并排序
unique_keys = df['key'].unique().compute().sort_values()

# 2. 生成分区边界(Dask的divisions是左闭右开区间,所以需要在最后加一个超出最大key的边界)
if len(unique_keys) > 0:
    # 对于十六进制字符串,生成一个比最大key大的字符串作为结尾边界
    max_key = unique_keys[-1]
    next_key = f"{int(max_key, 16)+1:02X}"
    divisions = list(unique_keys) + [next_key]
else:
    divisions = []

# 3. 设置索引并按指定分区边界重新分区,shuffle='disk'避免内存溢出
df = df.set_index('key', divisions=divisions, shuffle='disk')

# 4. 输出Parquet
dd.to_parquet(df, './output.parquet', write_index=False)

方案2:按key分组保存(适合Hive-style分区查询)

如果你希望输出的Parquet是Hive风格的分区结构(每个key对应一个子文件夹,比如key=00/、key=FF/),而且每个子文件夹里的文件只属于该key,可以用groupby的方式实现:

import dask.dataframe as dd

# 读取和生成key列的逻辑和上面一致
df = dd.read_csv(
    "/Users/ecerulm/Downloads/test/**/*.txt.gz",
    include_path_column=True,
    sep="\t",
    compression='gzip',
    blocksize=None
)

df['basename'] = df.path.str.rpartition('/')[2]
df['MAC'] = df.basename.str.partition('.')[2]
df['MAC'] = df.MAC.str.partition('.')[0]
df['key'] = df.MAC.str[-2:]

# 定义每个分组的保存函数
def save_single_key_group(group):
    # 获取当前组的key值
    current_key = group['key'].iloc[0]
    # 保存到对应子文件夹,支持自动分块(如果组内数据过大)
    group.to_parquet(f'./output_parquet/key={current_key}')

# 按key分组并执行保存,meta=object是告诉Dask返回的元数据类型
df.groupby('key').apply(save_single_key_group, meta=object).compute()

这种方式生成的结构更适合Spark、Pandas等其他工具按key快速查询,每个子文件夹下的所有文件都只包含该key的数据。

为什么原来的方法不生效?

你之前用的df.shuffle(['key']) + repartition(partition_size='100MB'):

  • shuffle(['key'])只是做哈希分区,无法保证绝对的纯净分区(存在哈希碰撞的可能)
  • repartition(partition_size='100MB')会再按大小拆分分区,可能把同一个key的数据拆到多个分区,也可能把不同key的小数据合并到一个分区

上面两种方案都能彻底解决你的需求,选哪种取决于你后续的使用场景~

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 20:22:48