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

如何预缓存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

  1. Broadcast the small lookup table to all workers: Instead of turning your small right DataFrame into a Dask-cuDF object, we’ll send the full in-memory cuDF DataFrame directly to every worker and cache it there using client.scatter() with the broadcast=True flag.
  2. 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_right across 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:27:05