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

Django大规模路由生成的故障处理与断点续传方案咨询

Django + Celery + Redis 实现大规模路由生成的断点续传方案

针对4.3万地点生成55亿条路由的场景,原方案故障后需从头重跑的问题,下面给出基于Celery+Redis的断点续传实现,包含模型调整、任务拆分和进度追踪的完整方案:

一、模型调整

新增进度追踪模型记录处理状态,同时给Route加唯一约束避免重复创建:

from django.db import models
from django.db.models import UniqueConstraint

class Location(models.Model):
    zip_code = models.CharField(max_length=5)

class Route(models.Model):
    origin = models.ForeignKey(
        Location, related_name="route_origin", on_delete=models.CASCADE
    )
    destination = models.ForeignKey(
        Location, related_name="route_destination", on_delete=models.CASCADE
    )

    VAN, REEFER, FLATBED = "V", "R", "F"
    equipment_type = models.CharField(
        max_length=50,
        choices=(
            (VAN, "Van"),
            (REEFER, "Reefer"),
            (FLATBED, "Flatbed"),
        ),
    )

    class Meta:
        # 唯一约束,确保同一起点、终点、设备类型的路由只存一次
        constraints = [
            UniqueConstraint(fields=['origin', 'destination', 'equipment_type'], name='unique_route')
        ]

# 单例进度模型,记录当前处理到的Origin ID
class RouteGenerationProgress(models.Model):
    last_processed_origin_id = models.IntegerField(default=0)
    is_completed = models.BooleanField(default=False)

    class Meta:
        verbose_name_plural = "路由生成进度"

    @classmethod
    def get_singleton(cls):
        # 确保只有一条进度记录
        obj, created = cls.objects.get_or_create(pk=1)
        return obj

二、Celery任务实现

把生成任务按Origin拆分,每次处理单个(或多个)Origin对应的所有路由,处理完更新进度。任务中断后,重启时自动从上次进度继续:

from celery import shared_task
from .models import Route, Location, RouteGenerationProgress
from django.db import transaction

@shared_task(bind=True, max_retries=3)
def generate_routes_batch(self):
    progress = RouteGenerationProgress.get_singleton()
    if progress.is_completed:
        return "路由生成已完成"
    
    # 获取未处理的Origin(ID大于上次处理的ID)
    unprocessed_origins = Location.objects.filter(id__gt=progress.last_processed_origin_id).order_by('id')
    if not unprocessed_origins.exists():
        progress.is_completed = True
        progress.save()
        return "所有路由生成完成"
    
    # 单次处理的Origin数量,根据服务器性能调整
    batch_size = 1
    current_origins = unprocessed_origins[:batch_size]

    for origin in current_origins:
        # 获取所有终点(如果允许起点=终点,去掉exclude)
        destinations = Location.objects.exclude(id=origin.id)
        equipment_types = [Route.VAN, Route.REEFER, Route.FLATBED]
        
        # 批量构造路由对象
        route_objects = []
        # 如果终点数量过大,可再拆分终点批次,避免内存溢出
        dest_batch_size = 1000
        for dest_batch in [destinations[i:i+dest_batch_size] for i in range(0, destinations.count(), dest_batch_size)]:
            for dest in dest_batch:
                for eq_type in equipment_types:
                    route_objects.append(Route(
                        origin=origin,
                        destination=dest,
                        equipment_type=eq_type
                    ))
        
        # 批量插入,ignore_conflicts跳过已存在的路由
        with transaction.atomic():
            Route.objects.bulk_create(route_objects, ignore_conflicts=True, batch_size=1000)
        
        # 更新进度到当前Origin的ID
        progress.last_processed_origin_id = origin.id
        progress.save()
    
    # 递归触发下一批任务
    generate_routes_batch.delay()
    return f"已完成Origin ID {progress.last_processed_origin_id} 的路由生成,继续下一批"

三、启动与断点续传

启动任务直接调用即可:

from .tasks import generate_routes_batch

# 启动路由生成
generate_routes_batch.delay()

如果任务中途中断(服务器重启、Celery异常),只需再次调用generate_routes_batch.delay(),任务会自动读取RouteGenerationProgress中的进度,从上次未完成的Origin开始继续生成,不会重复创建已存在的路由。

四、核心优化说明

  • 批量插入:用bulk_create替代循环create,配合batch_size控制插入量,比单条插入效率提升数十倍。
  • 事务保障:每个Origin的处理放在事务中,确保进度和数据一致,避免部分生成的情况。
  • 唯一约束:配合ignore_conflicts=True,即使重复触发任务也不会生成重复数据。
  • 批次拆分:拆分Origin和Destination的处理批次,避免内存占用过高,适配大规模数据场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 22:44:59