无法用dd.read_parquet时,如何指定Dask读取Parquet的分区数?
dd.from_delayed Great question! Since you’re constructing your Dask DataFrame using delayed tasks and dd.from_delayed (instead of the standard dd.read_parquet), you’ve got two solid options to control the number of final partitions:
1. Batch Row Groups During Task Creation (Most Efficient)
Instead of creating one delayed task per row group, bundle multiple row groups into a single delayed task. Each bundled task will become one partition in your final Dask DataFrame, letting you set the exact number of partitions upfront.
Here’s how to implement it:
import pandas as pd from dask import delayed, dataframe as dd # Define your target number of partitions target_partitions = 10 # Adjust this to your needs row_groups = list(pf.row_groups) # Calculate how many row groups to include per batch (round up to avoid missing groups) batch_size = (len(row_groups) + target_partitions - 1) // target_partitions # Create batches of row groups and define a delayed function to read each batch batched_tasks = [] for i in range(0, len(row_groups), batch_size): current_batch = row_groups[i:i+batch_size] @delayed def read_batch(rg_batch, columns, cats): # Read all row groups in the batch and concatenate into a single DataFrame batch_dfs = [pf.read_row_group_file(rg, columns, cats) for rg in rg_batch] return pd.concat(batch_dfs, ignore_index=True) batched_tasks.append(read_batch(current_batch, pf.columns, pf.cats)) # Build the Dask DataFrame from batched tasks df = dd.from_delayed(batched_tasks)
This approach avoids extra shuffle operations because you’re controlling partition count during the initial data loading phase. For best performance, you can even adjust batch sizes based on the size of each row group (if you have that info) to ensure balanced partition sizes.
2. Repartition After Creating the Dask DataFrame
If you already have your initial Dask DataFrame and just need to adjust partitions, use the repartition method. This is simpler but may trigger a shuffle of data across workers, which can add overhead for large datasets.
Example:
# After creating df with dd.from_delayed(dfs) df = df.repartition(npartitions=10)
Use this option if you need a quick fix or if you’re unsure about optimal batch sizes during the initial load.
内容的提问来源于stack exchange,提问作者j-bennet

