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

django-dramatiq流水线仅执行首个阶段问题排查

Django-Dramatiq 多阶段流水线仅执行首个Actor问题

我尝试使用django-dramatiq,通过pipeline(<stages>).run()方法运行由多个dramatiq Actor组成的多阶段流水线,但仅首个阶段/Actor会执行,其余阶段无任何执行动作。

测试用Actor定义

import dramatiq

@dramatiq.actor
def fake_extract(process_pk, *args, **kwargs):
    print(f"fake_extract: Process PK= {process_pk} Running extract on {kwargs['fits_file']}")

@dramatiq.actor
def fake_astromfit(process_pk, *args, **kwargs):
    print(f"fake_astromfit: Process PK= {process_pk} Astrometric fit on {kwargs['ldac_catalog']}, updating {kwargs['fits_file']}")

@dramatiq.actor
def fake_zeropoint(process_pk, *args, **kwargs):
    print(f"fake_zeropoint: Process PK= {process_pk} ZP determination on {kwargs['ldac_catalog']} with {kwargs['desired_catalog']} ref catalog")

流水线构建代码

import os
from dramatiq import pipeline
from test_dramatiq.dramatiq_tests import fake_extract, fake_astromfit, fake_zeropoint


fits_filepath = '/foo/bar.fits'
fits_file = os.path.basename(fits_filepath)

steps = [{
            'name'   : 'proc-extract',
            'runner' : fake_extract,
            'inputs' : {'fits_file':fits_filepath,
                       'datadir': os.path.join(dataroot, temp_dir)}
        },
        {
            'name'   : 'proc-astromfit',
            'runner' : fake_astromfit,
            'inputs' : {'fits_file' : fits_filepath,
                        'ldac_catalog' : os.path.join(dataroot, temp_dir, fits_file.replace('e91.fits', 'e91_ldac.fits')),
                        'datadir' : os.path.join(dataroot, temp_dir)
                        }
        },
        {
            'name'   : 'proc-zeropoint',
            'runner' : fake_zeropoint,
            'inputs' : {'ldac_catalog' : os.path.join(dataroot, temp_dir, fits_file.replace('e91.fits', 'e92_ldac.fits')),
                        'datadir' : os.path.join(dataroot, temp_dir),
                        'desired_catalog' : 'PS1'
                        }
        }]

pipes = []
for step_num, step in enumerate(steps):
    inputs = step['inputs']
    print(f"  Performing pipeline step {step['name']}")
    pk = 1234+step_num
    pipes.append(step['runner'].message_with_options(args=[pk,], kwargs=inputs, pipe_ignore=True))
pipeline(pipes).run()

运行结果对比

  • 常规Dramatiq(IPython中)运行正常:
fake_extract: Process PK= 1234 Running extract on /foo/bar.fits
fake_astromfit: Process PK= 1235 Astrometric fit on /foo/Temp_cvc2/bar.fits, updating /foo/bar.fits
fake_zeropoint: Process PK= 1236 ZP determination on /foo/Temp_cvc2/bar.fits with PS1 ref catalog
  • Django-Dramatiq环境下(python manage.py rundramatiq运行)仅执行首个阶段:
fake_extract: Process PK= 1234 Running extract on /foo/bar.fits

请问这是什么原因,为何django-dramatiq下多阶段流水线无法正常工作?


问题排查与解决方法

1. 移除pipe_ignore=True参数

pipe_ignore=True会让Dramatiq忽略消息的管道上下文,直接切断流水线阶段间的触发逻辑。构建流水线消息时不需要添加该参数,流水线本身会处理阶段依赖传递:

# 修改前
pipes.append(step['runner'].message_with_options(args=[pk,], kwargs=inputs, pipe_ignore=True))

# 修改后
pipes.append(step['runner'].message_with_options(args=[pk,], kwargs=inputs))

2. 确认Actor被Django-Dramatiq正确加载

确保所有Actor都能被django-dramatiq扫描到:

  • 检查settings.py中INSTALLED_APPS包含Actor所在的应用
  • 确保DRAMATIQ_AUTODISCOVER配置为True,或在应用的apps.py中手动导入Actor模块,避免仅首个Actor被加载

3. 验证Broker配置一致性

检查Django环境下的Broker(如Redis、RabbitMQ)配置和常规Dramatiq使用的完全一致。如果配置不同,后续阶段的消息可能发送到了另一个Broker实例,导致Worker无法接收。示例配置:

DRAMATIQ_BROKER = {
    "BROKER": "dramatiq.brokers.redis.RedisBroker",
    "OPTIONS": {
        "url": "redis://localhost:6379/0",
    },
    "MIDDLEWARE": [
        "dramatiq.middleware.Prometheus",
        "dramatiq.middleware.AgeLimit",
        "dramatiq.middleware.TimeLimit",
        "dramatiq.middleware.Callbacks",
        "dramatiq.middleware.Retries",
        "django_dramatiq.middleware.DjangoMiddleware",
    ]
}

4. 排查隐藏异常

django-dramatiq的DjangoMiddleware会处理Django上下文,若后续Actor执行时出现未捕获异常(如数据库连接、路径不存在等),可能会静默失败。建议在Actor中添加异常捕获与日志输出:

import logging
import dramatiq

logger = logging.getLogger(__name__)

@dramatiq.actor
def fake_astromfit(process_pk, *args, **kwargs):
    try:
        print(f"fake_astromfit: Process PK= {process_pk} Astrometric fit on {kwargs['ldac_catalog']}, updating {kwargs['fits_file']}")
    except Exception as e:
        logger.error(f"fake_astromfit failed for PK {process_pk}: {str(e)}", exc_info=True)
        raise

查看Django日志文件,确认是否有异常信息。


内容的提问来源于stack exchange,提问作者astrosnapper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 21:48:31