Dask处理1.2GB数据性能优化问询:并行计算未达预期
Let's break down what's going on here and fix that frustrating performance bottleneck you're seeing with Dask. The core issue right now is that your current Dask implementation isn't actually leveraging parallel computing—you're processing files one at a time with generators, which defeats the whole purpose of using Dask! Here's how to turn things around:
1. Stop Using Generators for Sequential Processing
Your current generator-based code (dataframes = (pd.read_csv(fname) for fname in filenames) followed by next(help) calls) processes files serially, one after another. Dask can't parallelize this because you're manually controlling the iteration. Instead, let Dask handle the bulk file loading natively:
import dask.dataframe as dd from dask.distributed import Client # Start a Dask client to enable parallelism and monitoring client = Client() # This launches a local cluster with your 8 cores # Let Dask read all files in parallel, auto-partitioning the data ddf = dd.read_csv(filenames) # Chain your data cleaning/transformation functions directly on the Dask DataFrame ddf = ddf.pipe(remove_unnessasey_data, l1) ddf = ddf.pipe(change_structure)
2. Fix Your Data Transformation Logic
Looking at your initial data processing code, there are a few critical mistakes that hurt both correctness and performance:
# ❌ Your original code has unassigned operations and drops the DataFrame structure df['DateTime']=dd.to_datetime(df['DateTime']) df['KWH/hh (per half hour) '].astype(float) # This line does nothing—no assignment! df=df['KWH/hh (per half hour) '].fillna(0) # Turns DataFrame into a Series (loses DateTime!) df=df.set_index(df['DateTime'], npartitions='auto') # Will fail because Series has no 'DateTime' column
Correct the Transformation Chain
Keep your data in a Dask DataFrame (don't collapse to a Series) and assign operations properly:
# ✅ Properly transform columns while preserving structure ddf = ddf.assign( # Parse DateTime as a datetime type DateTime=lambda df: dd.to_datetime(df['DateTime']), # Clean the KWH column: cast to float, fill NaNs with 0 KWH=lambda df: df['KWH/hh (per half hour) '].astype(float).fillna(0) ) # Set DateTime as the index for efficient resampling ddf = ddf.set_index('DateTime', npartitions=16) # Use 1-2x your core count (8 cores → 16 partitions) # Resample and compute the daily sum daily_kwh_sum = ddf['KWH'].resample('D').sum().compute()
- Why this works: Dask can parallelize operations across partitions when you keep the DataFrame structure. Assigning columns explicitly ensures all transformations are tracked and executed in parallel.
3. Tune Partition Sizing
The npartitions='auto' setting can sometimes create partitions that are too large (leading to slow processing per partition) or too small (leading to excessive scheduling overhead). For 1.2GB of data and 8 cores:
- Aim for partitions between 100-200MB (this is a sweet spot for Dask)
- Set
npartitions=16(twice your core count) to balance parallelism and overhead - You can check partition sizes with
ddf.memory_usage(deep=True).compute()to adjust as needed
4. Use Dask's Dashboard to Diagnose Bottlenecks
When you start the Dask Client(), open http://localhost:8787 in your browser. This dashboard will show you:
- Which tasks are taking the longest
- If any workers are idle
- Memory usage across partitions
This is invaluable for spotting issues like uneven partition sizes or slow transformation functions.
5. Compare to Pandas Chunks (and Why Dask Should Be Faster)
Your pandas chunks approach is faster right now because it's using efficient single-process chunking, but Dask should outperform it when configured correctly. The key difference is Dask can use all 8 of your cores in parallel, while pandas chunks are single-threaded. Once you fix the parallelism issue, Dask will process your 1.2GB data in a fraction of the time (way under an hour!).
内容的提问来源于stack exchange,提问作者Elham Eskandarnia

