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

Django+Celery异步进程日志写入数据库字段丢失问题排查与修复

Django + Celery 异步日志并发写入问题分析与修复

问题背景

基于Django和Celery的大文件分步处理项目中,每个文件对应一个Sequencing对象。处理流程会创建/复用samples、variants等对象,需求是将各区域的对象创建/复用数量写入Sequencing的metadata JSON字段。但异步执行时日志仅部分写入,有时只显示单个区域日志,有时显示部分区域,无固定规律,推测是load_variants_by_contig异步调用log方法的save操作时发生并发冲突。

冲突核心原因

多个Celery异步任务同时读取同一个Sequencing对象的metadata:

  1. 任务A读取当前metadata,添加新日志条目
  2. 同时任务B也读取同一版本的metadata,添加另一个日志条目
  3. 若任务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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 10:07:04