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
相关产品推荐
相关产品推荐

