如何预缓存Dask DataFrame至所有Worker,优化不平衡Merge通信效率?
Optimizing Left Joins with Small Lookup Tables in Dask + cuDF
When you're working with a large Dask-cuDF DataFrame paired with a small lookup table (like your right DataFrame here), pre-caching that lookup table on every worker is a huge win for cutting down network overhead and unnecessary data shuffling. Let me walk you through exactly how to implement this using client.scatter():
Step-by-Step Breakdown
- Broadcast the small lookup table to all workers: Instead of turning your small
rightDataFrame into a Dask-cuDF object, we’ll send the full in-memory cuDF DataFrame directly to every worker and cache it there usingclient.scatter()with thebroadcast=Trueflag. - Use the cached table in
map_partitions: When merging each partition of the large left DataFrame, we’ll reference the locally cached lookup table instead of a Dask DataFrame—so each worker can run the merge without pulling data across the network.
Modified Working Code
import pandas as pd import cudf import dask_cudf import numpy as np from dask.distributed import Client from dask_cuda import LocalCUDACluster # Create a local CUDA cluster cluster = LocalCUDACluster() client = Client(cluster) np.random.seed(12) nrows_left = 1000000 nrows_right = 1000 # Create local cuDF DataFrames left = cudf.DataFrame({'a': np.random.randint(0,nrows_right,nrows_left), 'left_value':np.arange(nrows_left)}) right = cudf.DataFrame({'a': np.arange(nrows_right), 'lookup_val': np.random.randint(0,1000,nrows_right)}) print(left.shape, right.shape) # (1000000, 2) (1000, 2) # Convert large left DataFrame to Dask-cuDF (split into many partitions) ddf_left = dask_cudf.from_cudf(left, npartitions=500) # Broadcast the small right lookup table to ALL workers (cached locally on each) scattered_right = client.scatter(right, broadcast=True) def dask_merge(L, R): # R is now the full cached cuDF DataFrame on each worker return L.merge(R, how='left', on='a') # Use the scattered/cached right table in map_partitions result = ddf_left.map_partitions(dask_merge, R=scattered_right).compute() result.head() # <cudf.DataFrame ncols=3 nrows=5 > # a left_value lookup_val # 0 219 1952 822 # 1 873 1953 844 # 2 908 1954 142 # 3 290 1955 810 # 4 863 1956 910
Why This Is More Efficient
- No repeated network transfers: Without broadcasting, Dask might repeatedly pull partitions of
ddf_rightacross the network for each merge operation. Broadcasting sends the table once to every worker and keeps it cached. - No data shuffling: Since every worker has the full lookup table, we don’t need to rearrange partitions of either DataFrame to align them—each left partition merges directly with the local copy of the lookup table.
- Lower memory bloat: Each worker holds just one copy of the small lookup table, instead of multiple copies from repeated partition pulls.
A quick note: this trick only makes sense for small lookup tables. If your right DataFrame is large enough that storing a copy on every worker would strain memory, stick with regular Dask merge operations instead.
内容的提问来源于stack exchange,提问作者Nick Becker
相关产品推荐
相关产品推荐

