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

