Django+Celery异步进程日志写入数据库字段丢失问题排查与修复
Django + Celery 异步日志并发写入问题分析与修复
问题背景
基于Django和Celery的大文件分步处理项目中,每个文件对应一个Sequencing对象。处理流程会创建/复用samples、variants等对象,需求是将各区域的对象创建/复用数量写入Sequencing的metadata JSON字段。但异步执行时日志仅部分写入,有时只显示单个区域日志,有时显示部分区域,无固定规律,推测是load_variants_by_contig异步调用log方法的save操作时发生并发冲突。
冲突核心原因
多个Celery异步任务同时读取同一个Sequencing对象的metadata:
- 任务A读取当前
metadata,添加新日志条目 - 同时任务B也读取同一版本的
metadata,添加另一个日志条目 - 若任务A先执行
save,任务B后续save时会用自己之前读取的旧metadata覆盖任务A的修改,导致部分日志丢失。
调试验证方法
- 添加任务级日志:在
log方法中打印任务ID、操作前后的日志长度,追踪并发覆盖情况:import os from django.utils import timezone def log(self, level: str, msg: str | list, *args, **kwargs) -> None: task_id = os.getenv('CELERY_TASK_ID', 'local') before_len = len(self.metadata.get('log', [])) # 原有日志添加逻辑... self.save() # 重新读取最新数据验证 after_obj = Sequencing.objects.get(pk=self.pk) after_len = len(after_obj.metadata.get('log', [])) print(f"Task {task_id}: Before save log len {before_len}, after {after_len}") - 行级锁验证:在任务中获取
Sequencing对象时使用select_for_update()加行级锁,若日志恢复正常,则确认是并发冲突:# tasks.py中修改任务 sequencing_obj = Sequencing.objects.select_for_update().get(pk=sequencing_pk) - 查看Celery执行日志:开启Celery详细日志(
celery -A your_project worker -l info),观察任务执行顺序和save操作的时间点,确认并发写入的发生时机。
修复方案
方案1:行级锁+字段更新优化(快速修复)
通过数据库行级锁确保同一时间只有一个任务修改Sequencing对象,同时优化log方法只更新必要字段:
# tasks.py 异步任务修改 @shared_task(bind=True, queue="computation") def sequencing_load_variants_by_contig(self, sequencing_pk: int, contig: str) -> None: from sequencings.models import Sequencing from django.db import transaction with transaction.atomic(): # 加行级锁,阻塞其他任务读取该对象,直到当前事务结束 sequencing_obj = Sequencing.objects.select_for_update().get(pk=sequencing_pk) sequencing_obj.load_variants_by_contig(contig=contig) # models.py log方法优化 def log(self, level: str, msg: str | list, *args, **kwargs) -> None: # 刷新获取最新的metadata,避免使用缓存的旧数据 self.refresh_from_db(fields=['metadata']) if "log" not in self.metadata: self.metadata["log"] = [] self.metadata["log"].append( { "date": datetime.datetime.utcnow().isoformat(), "level": level, "msg": msg, } ) # 仅更新metadata字段,提升性能 self.save(update_fields=['metadata'])
方案2:独立日志表(推荐,适合长期扩展)
将日志从JSON字段拆分到独立表,彻底避免并发覆盖问题,同时方便后续查询和统计:
# models.py 新增日志模型 class SequencingLog(models.Model): SEVERITY_CHOICES = [ ('info', 'Info'), ('warning', 'Warning'), ('error', 'Error'), ] sequencing = models.ForeignKey(Sequencing, on_delete=models.CASCADE, related_name='logs') date = models.DateTimeField(auto_now_add=True) level = models.CharField(max_length=20, choices=SEVERITY_CHOICES) msg = models.JSONField() # 兼容字符串或列表格式 # 修改Sequencing的log方法 def log(self, level: str, msg: str | list, *args, **kwargs) -> None: SequencingLog.objects.create( sequencing=self, level=level, msg=msg )
后续查看日志可直接通过sequencing_obj.logs.all()获取,无需操作JSON字段。
方案3:Redis暂存+批量写入(适合高并发场景)
用Redis暂存日志条目,所有区域处理完成后批量写入数据库,减少并发写入次数:
# models.py log方法修改 import redis import json from django.conf import settings from django.utils import timezone def log(self, level: str, msg: str | list, *args, **kwargs) -> None: r = redis.Redis(host=settings.REDIS_HOST, port=settings.REDIS_PORT, db=0) log_entry = { "date": timezone.now().isoformat(), "level": level, "msg": msg, } # 将日志推入Redis队列 r.rpush(f"sequencing_log:{self.pk}", json.dumps(log_entry)) # tasks.py 新增批量写入任务 @shared_task(bind=True, queue="computation") def batch_write_sequencing_logs(self, sequencing_pk: int) -> None: from sequencings.models import Sequencing import redis import json from django.conf import settings r = redis.Redis(host=settings.REDIS_HOST, port=settings.REDIS_PORT, db=0) sequencing_obj = Sequencing.objects.get(pk=sequencing_pk) log_key = f"sequencing_log:{sequencing_pk}" # 从Redis取出所有日志 logs = [] while True: entry = r.lpop(log_key) if not entry: break logs.append(json.loads(entry)) if logs: if "log" not in sequencing_obj.metadata: sequencing_obj.metadata["log"] = [] sequencing_obj.metadata["log"].extend(logs) sequencing_obj.save(update_fields=['metadata']) # 修改load_variants方法,用chord确保所有区域任务完成后执行批量写入 def load_variants(self) -> None: contigs = [...] # 创建所有区域处理任务的组 tasks_group = group( sequencing_load_variants_by_contig.si(self.pk, contig=contig) for contig in contigs ) # 组任务全部完成后触发批量写日志 chord(tasks_group)(batch_write_sequencing_logs.si(self.pk))
内容的提问来源于stack exchange,提问作者Max B 0815
相关产品推荐
相关产品推荐

