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

如何配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:10:58