如何将Celery任务爬取数据存临时库后通过Django ORM比对合并至主库
Django ORM 增量合并爬虫新旧数据方案
核心思路是基于唯一标识做增量更新,无需清空全表,仅处理变化的数据,资源消耗比全量覆盖低90%以上。
前置准备
首先确定黑客松数据的唯一识别字段,比如官方活动ID+数据源组合、活动官网唯一链接等,保证该组合不会出现重复,作为新旧数据匹配的依据。
方案1:update_or_create 逐行处理(适合中小数据量)
如果每日爬取的黑客松数据量在1万条以内,直接用Django ORM原生的update_or_create方法即可,写法简单不易出错:
from django.db import transaction from .models import Hackathon def merge_hackathon_data(scraped_items): # 事务包裹避免逐行提交,大幅提升写入性能 with transaction.atomic(): for item in scraped_items: Hackathon.objects.update_or_create( # 唯一匹配条件,可根据实际业务调整 official_id=item["official_id"], source=item["source"], # 匹配到旧记录则更新以下字段,未匹配则新建整条记录 defaults={ "name": item["name"], "start_time": item["start_time"], "end_time": item["end_time"], "location": item["location"], "register_link": item["register_link"], "tags": item["tags"] # 其他需要更新的字段 } )
该方案优势:原生ORM支持,逻辑简单,自带原子性,不会出现部分更新成功部分失败的问题。
方案2:批量比对+批量操作(适合大数据量)
如果每日爬取数据超过5万条,逐行处理性能较差,可采用先批量比对再批量写入的方案,性能比逐行处理高3~10倍:
from django.db import transaction from .models import Hackathon def merge_hackathon_data(scraped_items): # 构造唯一键到爬取数据的映射 item_key_map = {} unique_keys = [] for item in scraped_items: key = (item["official_id"], item["source"]) item_key_map[key] = item unique_keys.append(key) # 批量查询已存在的旧记录 existing_records = Hackathon.objects.filter( official_id__in=[k[0] for k in unique_keys], source__in=[k[1] for k in unique_keys] ).only("id", "official_id", "source") # 拆分待更新、待新建数据集 existing_key_to_id = {(r.official_id, r.source): r.id for r in existing_records} to_create = [] to_update = [] for key, item in item_key_map.items(): if key in existing_key_to_id: item["id"] = existing_key_to_id[key] to_update.append(Hackathon(**item)) else: to_create.append(Hackathon(**item)) with transaction.atomic(): # 批量插入新记录 if to_create: Hackathon.objects.bulk_create(to_create, batch_size=1000) # 批量更新旧记录,指定更新字段提升性能 if to_update: Hackathon.objects.bulk_update( to_update, fields=["name", "start_time", "end_time", "location", "register_link", "tags"], batch_size=1000 )
可选:过期数据处理
如果需要标记已经下架的旧黑客松活动,可以给模型加is_active布尔字段,合并完成后批量标记不在新爬数据集内的旧记录为非活跃:
# 合并完成后执行,标记过期活动 Hackathon.objects.exclude( official_id__in=[k[0] for k in unique_keys], source__in=[k[1] for k in unique_keys] ).update(is_active=False)
Celery任务优化建议
- 可以将爬取、合并拆分为两个独立子任务,避免单任务执行时间过长超时
- 数据量极大的情况下可以按数据源分片处理,多个子任务并行执行合并逻辑
- 可以设置任务超时时间和失败重试机制,避免网络波动导致任务失败
内容的提问来源于stack exchange,提问作者A.K.M. Adib
相关产品推荐
相关产品推荐

