如何配置Dask-CUDA实现多GPU间Dask-CUDF工作负载均衡?
Dask-CUDF多GPU负载分配问题解决指南
问题背景
拥有两块NVIDIA GeForce RTX 3090 GPU,需要用dask_cudf将数据处理任务均匀分配到多GPU,但当前工作负载仅由单块GPU处理。以下是简化版脚本:
import os import json import cudf import dask_cudf import time from concurrent.futures import ThreadPoolExecutor from dask.distributed import Client from dask_cuda import LocalCUDACluster import socket def load_data(): # Placeholder for data loading logic return {"dummy": dask_cudf.DataFrame()} def handle_client(client_socket, datasets): try: data = client_socket.recv(4096).decode('utf-8').strip() if data: session_id, words_json = data.split(' ', 1) words = json.loads(words_json) combined_df = dask_cudf.DataFrame() total_ile = 0 for word in words: if word in datasets: combined_df = dask_cudf.concat([combined_df, datasets[word]]) total_ile += datasets[word]['ile'].unique().sum().compute() if not combined_df.empty: id_counts = combined_df['id'].value_counts().reset_index() id_counts.columns = ['id', 'count'] combined_df = combined_df.merge(id_counts, on='id', how='left') combined_df['score'] = combined_df['score'] + (combined_df['count'] - 1) * 333 grouped_df = combined_df.groupby('id').agg({'score': 'sum', 'domain_id': 'first'}).reset_index() top_results = grouped_df.nlargest(1000, 'score').compute().to_pandas().to_dict('records') else: top_results = [] response = json.dumps({ 'session_id': session_id, 'data': top_results, 'response_time': time.time() - start_time, 'ile': total_ile }) client_socket.sendall(response.encode('utf-8')) client_socket.sendall(b"END\n") except json.JSONDecodeError: client_socket.sendall(b"END\n") finally: client_socket.close() def initialize_cluster(): cluster = LocalCUDACluster( protocol='tcp', dashboard_address=':8787', CUDA_VISIBLE_DEVICES="0,1", rmm_pool_size="16GB" ) client = Client(cluster) return client def main(): client = initialize_cluster() datasets = load_data() server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) server_socket.bind(("localhost", 9999)) server_socket.listen(5) try: with ThreadPoolExecutor() as executor: while True: client_socket, addr = server_socket.accept() executor.submit(handle_client, client_socket, datasets) finally: server_socket.close() if __name__ == "__main__": main()
用户疑问
- 如何配置dask_cuda及脚本,实现工作负载在两块GPU间均匀分配?
- 是否存在遗漏的特定设置或配置?
解决方案与关键修改
一、修正集群初始化配置
确保LocalCUDACluster明确启动对应GPU数量的worker,避免默认逻辑的不确定性:
def initialize_cluster(): cluster = LocalCUDACluster( protocol='tcp', dashboard_address=':8787', CUDA_VISIBLE_DEVICES="0,1", rmm_pool_size="16GB", n_workers=2, # 显式指定worker数量与GPU数一致 threads_per_worker=1 # 每个worker用单线程,避免CPU线程竞争 ) client = Client(cluster) print(f"Dask Dashboard: {client.dashboard_link}") # 输出监控链接 return client
二、保证数据分区与集群分布
Dask-CUDF的多GPU调度依赖数据分区,必须让数据拥有至少与GPU数量相等的分区,并持久化到集群:
def load_data(client): # 传入client用于数据集群分布 # 替换为真实数据加载逻辑,示例用模拟数据演示分区 dummy_cudf = cudf.DataFrame({ 'id': range(100000), 'ile': [1]*100000, 'score': range(100000), 'domain_id': [0]*100000 }) # 生成2个分区,对应两块GPU dummy_dask_df = dask_cudf.from_cudf(dummy_cudf, npartitions=2) # 将数据持久化到集群,让分区分布到不同GPU worker dummy_dask_df = dummy_dask_df.persist() return {"dummy": dummy_dask_df}
同时修改main函数的调用:
def main(): client = initialize_cluster() datasets = load_data(client) # 传入client # ... 后续socket代码不变
三、修正任务执行逻辑,通过Dask集群调度计算
避免在本地线程直接调用compute(),改用Dask Client提交任务到集群,让任务自动分配到多GPU:
def handle_client(client_socket, datasets, client): # 传入client start_time = time.time() # 修复原代码中start_time未定义问题 try: data = client_socket.recv(4096).decode('utf-8').strip() if data: session_id, words_json = data.split(' ', 1) words = json.loads(words_json) combined_df = dask_cudf.DataFrame() total_ile = 0 for word in words: if word in datasets: combined_df = dask_cudf.concat([combined_df, datasets[word]]) # 异步提交ile计算任务到集群 ile_future = client.submit(lambda df: df['ile'].unique().sum().compute(), datasets[word]) total_ile += ile_future.result() if not combined_df.empty: # 保持Dask操作的惰性,最后统一提交计算 id_counts = combined_df['id'].value_counts().reset_index() id_counts.columns = ['id', 'count'] combined_df = combined_df.merge(id_counts, on='id', how='left') combined_df['score'] = combined_df['score'] + (combined_df['count'] - 1) * 333 grouped_df = combined_df.groupby('id').agg({'score': 'sum', 'domain_id': 'first'}).reset_index() top_results_df = grouped_df.nlargest(1000, 'score') # 提交任务到集群,由Dask调度到多GPU执行 top_results_future = client.compute(top_results_df) top_results = top_results_future.result().to_pandas().to_dict('records') else: top_results = [] response = json.dumps({ 'session_id': session_id, 'data': top_results, 'response_time': time.time() - start_time, 'ile': total_ile }) client_socket.sendall(response.encode('utf-8')) client_socket.sendall(b"END\n") except json.JSONDecodeError: client_socket.sendall(b"END\n") finally: client_socket.close()
修改main函数中线程池的调用:
executor.submit(handle_client, client_socket, datasets, client)
关键注意事项与遗漏点排查
- 数据分区是核心:如果数据只有1个分区,所有计算都会绑定到单GPU,必须保证
npartitions >= GPU数量,可通过repartition(npartitions=2)重新分区。 - 必须持久化数据:未执行
persist()的Dask DataFrame会每次计算重新加载,无法固定到多GPU节点。 - 避免本地计算:直接在主线程/线程池线程调用
compute()会绕过Dask集群调度,导致任务仅在本地GPU执行。 - 监控集群状态:通过Dask Dashboard(默认
http://localhost:8787)查看任务分布,确认两块GPU均有任务执行。
内容的提问来源于stack exchange,提问作者allo allo
相关产品推荐
相关产品推荐

