如何按指定列值拆分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筛选。
步骤如下:
- 先计算所有唯一的
key值并排序,用来生成分区边界 - 设置索引时指定这些边界,确保每个分区只包含单个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
相关产品推荐
相关产品推荐

