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

Django与Dask集成优化:性能、最佳实践及进度查询改进

关于性能与最佳实践

注:以下问题的完整代码已在Github公开(链接已移除)

我正尝试实现一个简单的Django+Dask集成:一个视图启动长耗时任务,另一个视图可查询任务状态,后续计划让get_task_status(或其他Django视图函数)返回任务输出。

我用time.sleep(2)模拟长耗时任务,且需要看到任务整体状态为"running",为此在测试中也使用了time.sleep(),这感觉不太合理。

视图代码如下:

from uuid import uuid4
from django.http import JsonResponse
from dask.distributed import Client
import time

# Initialize Dask client
client = Client(n_workers=8, threads_per_worker=2)

NUM_FAKE_TASKS = 25

# Dictionary to store futures with task_id as key
task_futures = {}


def long_running_process(work_list):
    def task_function(task):
        time.sleep(2)
        return task

    futures = [client.submit(task_function, task) for task in work_list]
    return futures


def start_task(request):
    work_list = []

    for t in range(NUM_FAKE_TASKS):
        task_id = str(uuid4())  # Generate a unique ID for the task
        work_list.append(
            {"address": f"foo--{t}@example.com", "message": f"Mail task: {task_id}"}
        )

    futures = long_running_process(work_list)
    dask_task_id = futures[0].key  # Use the key of the first future as the task ID

    # Store the futures in the dictionary with task_id as key
    task_futures[dask_task_id] = futures

    return JsonResponse({"task_id": dask_task_id})


def get_task_status(request, task_id):
    futures = task_futures.get(task_id)

    if futures:
        if not all(future.done() for future in futures):
            progress = 0
            return JsonResponse({"status": "running", "progress": progress})
        else:
            results = client.gather(futures, asynchronous=False)

            # Calculate progress, based on futures that are 'done'
            progress = int((sum(future.done() for future in futures) / len(futures)) * 100)

            return JsonResponse(
                {
                    "task_id": task_id,
                    "status": "completed",
                    "progress": progress,
                    "results": results,
                }
            )
    else:
        return JsonResponse({"status": "error", "message": "Task not found"})

我编写的测试耗时约5.5秒:

from django.test import Client
from django.urls import reverse
import time


def test_immediate_response_with_dask():
    client = Client()
    response = client.post(reverse("start_task_dask"), data={"data": "foo"})
    assert response.status_code == 200
    assert "task_id" in response.json()

    task_id = response.json()["task_id"]
    response2 = client.get(reverse("get_task_status_dask", kwargs={"task_id": task_id}))
    assert response2.status_code == 200
    r2_status = response2.json()["status"]
    assert r2_status == "running"

    attempts = 0
    max_attempts = 8

    while attempts < max_attempts:
        time.sleep(1)
        try:
            response3 = client.get(
                reverse("get_task_status_dask", kwargs={"task_id": task_id})
            )
            assert response3.status_code == 200

            r3_status = response3.json()["status"]
            r3_progress = response3.json()["progress"]

            assert r3_progress >= 99
            assert r3_status == "completed"
            break  # Exit the loop if successful
        except Exception:
            attempts += 1
            if attempts == max_attempts:
                raise  # Raise the last exception if all attempts failed

我的问题是:实现该API是否有更高效的方式?若NUM_FAKE_TASKS = 10000该如何处理?当前实现是否存在资源浪费?

补充:如何查看进度百分比?

已知以下代码是阻塞的:

client.gather(futures, asynchronous=False)

但client.gather(futures, asynchronous=True)的表现也不符合预期。

是否可通过client.persist()或client.compute()实现增量进度查看?我无法持久化list类型的<class 'distributed.client.Future'>,且使用client.compute(futures)时进度会直接从0跳至100,不符合预期。


内容的提问来源于stack exchange,提问作者blong

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 03:24:55