在Django中使用multiprocessing.Pool批量更新百万级数据的实现方案
用multiprocessing.Pool批量处理百万级Django记录
直接遍历数百万条记录逐个处理效率极低,我们可以用multiprocessing.Pool分批次并行处理,每1000条为一个批次,以下是具体实现方案:
核心思路
- 按ID分段获取数据:避免一次性加载所有记录到内存,通过ID范围拆分批次,减少内存占用并提升查询效率
- 多进程并行处理:利用进程池同时处理多个批次,充分利用CPU资源
- 进程内独立初始化:每个进程单独建立数据库连接,避免多进程共享连接导致的问题
代码示例
1. 批次处理函数
import django django.setup() # 独立脚本需添加,Django管理命令可省略 from myapp.models import MyTable def process_batch(id_range): start_id, end_id = id_range # 按ID范围查询当前批次的记录 entries = MyTable.objects.filter(id__gte=start_id, id__lte=end_id) for entry in entries: do_some_magic(entry) # 包含entry.save()的业务逻辑 return f"完成批次处理:{start_id}-{end_id}"
2. 主逻辑(分批次+进程池调度)
from multiprocessing import Pool from django.db.models import Min, Max from myapp.models import MyTable def main(): batch_size = 1000 # 获取记录的最小/最大ID,用于生成批次范围 id_stats = MyTable.objects.aggregate(min_id=Min('id'), max_id=Max('id')) min_id, max_id = id_stats['min_id'], id_stats['max_id'] if not min_id or not max_id: print("无记录需要处理") return # 生成所有批次的ID区间 batches = [] current_start = min_id while current_start <= max_id: current_end = min(current_start + batch_size - 1, max_id) batches.append((current_start, current_end)) current_start = current_end + 1 # 初始化进程池(默认大小为CPU核心数) with Pool() as pool: # 用imap_unordered实时获取处理结果,也可改用map等待所有批次完成 for result in pool.imap_unordered(process_batch, batches): print(result) if __name__ == "__main__": main()
关键注意事项
- Django环境初始化:如果是独立运行的脚本,必须调用
django.setup();如果是自定义Django管理命令,无需额外处理 - 数据库连接隔离:每个进程会自动创建独立的数据库连接,不要在进程间共享数据库连接对象
- 业务函数可序列化:
do_some_magic函数及相关数据需支持序列化,因为多进程间会传递任务参数 - ID不连续的情况:即使记录ID不连续,按范围查询依然有效,只是部分批次的记录数可能少于1000
- 性能调优:可根据服务器CPU核心数手动设置进程池大小(如
Pool(processes=4)),避免进程过多导致资源竞争
内容的提问来源于stack exchange,提问作者vukojevicf
相关产品推荐
相关产品推荐

