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

Flask+BullMQ Worker问题:任务已加入队列但未执行

问题描述

使用Flask结合BullMQ处理报告生成后台任务,任务已成功加入队列,日志显示Worker初始化正常(输出“Worker is ready and listening for jobs.”),但报告生成任务始终未执行。

完整代码

from flask import Flask, request, jsonify
from bullmq import Queue, Worker
import asyncio
from Mytools import *
import signal
import sys
import logging
import threading

app = Flask(__name__)
queue = Queue("generate_reports_queue")

# Configure logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

# ========================== #
#  Report Generation Function
# ========================== #
async def generate_report_task(survey_id):
    try:
        logger.info(f"Starting report generation for survey ID: {survey_id}")
        reponses, questions = fetch_responses(survey_id)
        data_df = responses_json_to_df(reponses, "definition")
        survey_data = bdv_json_to_df(questions["definition"].iloc[0])

        output_path = f"output_{survey_id}.docx"

        decision_maker_with_report_streaming(
            data=data_df,
            df_bbl=survey_data,
            template="Template_.docx",
            bucket_name="reu-data",
            key=output_path
        )

        logger.info(f"Report generated successfully for survey ID: {survey_id}")
        return {"message": "Report generated successfully", "file_path": output_path}

    except Exception as e:
        logger.error(f"Error generating report for survey ID {survey_id}: {str(e)}")
        return {"error": str(e)}

# ========================== #
#  Flask API
# ========================== #
@app.route('/generate-report', methods=['POST'])
async def generate_report():
    data = request.get_json()
    survey_id = data.get("id")

    if not survey_id:
        return jsonify({"error": "ID missing"}), 400

    job = await queue.add("generate_report", {"survey_id": survey_id})
    logger.info(f"Added job to queue with ID: {job.id}")

    return jsonify({"message": "Task in progress", "task_id": job.id})

# ========================== #
#  BullMQ Worker
# ========================== #
async def process_job(job):
    survey_id = job.data.get("survey_id")
    logger.info(f"Processing job with survey ID: {survey_id}")
    return await generate_report_task(survey_id)

async def main_worker():
    try:
        logger.info("Initializing worker...")
        worker = Worker("generate_reports_queue", process_job)
        logger.info("Worker is ready and listening for jobs.")
        await asyncio.Future()  # Keep the worker running
    except Exception as e:
        logger.error(f"Worker error: {e}", exc_info=True)

def run_worker():
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    loop.run_until_complete(main_worker())

# Graceful shutdown handling
def shutdown(signal, frame):
    logger.info("Shutting down gracefully...")
    sys.exit(0)

signal.signal(signal.SIGINT, shutdown)
signal.signal(signal.SIGTERM, shutdown)

if __name__ == "__main__":
    worker_thread = threading.Thread(target=run_worker, daemon=True)
    worker_thread.start()
    app.run(debug=True, use_reloader=False)

预期行为

  • Worker拾取任务后调用generate_report_task函数并生成报告。

实际行为

  • 任务已成功加入队列。
  • 日志输出"Worker is ready and listening for jobs."。
  • generate_report_task函数从未被调用。

测试API日志

* Serving Flask app 'flask_api_with_bull2'
INFO:__main__:Initializing worker...
 * Debug mode: on
INFO:__main__:Worker is ready and listening for jobs.
INFO:werkzeug:WARNING: This is a development server. Do not use it in a production deployment. Use a production WSGI server instead.
 * Running on http://127.0.0.1:5000
INFO:werkzeug:Press CTRL+C to quit
INFO:__main__:Added job to queue with ID: 8
INFO:werkzeug:127.0.0.1 - - [10/Mar/2025 15:18:01] "POST /generate-report HTTP/1.1" 200 -

问题排查与解决方法

1. Worker未启动核心循环(最可能的原因)

代码仅初始化了Worker对象,但未调用BullMQ Worker的run()方法启动任务监听循环。当前用await asyncio.Future()保持线程运行,但Worker并未真正开始处理队列任务。

修复方法:修改main_worker函数:

async def main_worker():
    try:
        logger.info("Initializing worker...")
        worker = Worker("generate_reports_queue", process_job)
        logger.info("Worker is ready and listening for jobs.")
        await worker.run()  # 启动Worker任务处理循环
    except Exception as e:
        logger.error(f"Worker error: {e}", exc_info=True)

2. Redis连接异常

BullMQ依赖Redis作为消息中间件,默认连接本地Redis(localhost:6379)。若Redis未运行、端口/地址配置错误,队列和Worker将无法通信,任务会滞留在队列中。

排查步骤:

  • 检查本地Redis服务是否启动(执行redis-cli ping,返回PONG则正常)。
  • 显式指定Redis连接参数,确保Queue和Worker使用相同配置:
# 初始化Queue时指定连接
queue = Queue("generate_reports_queue", connection={"host": "localhost", "port": 6379})

# 初始化Worker时同样指定连接
worker = Worker("generate_reports_queue", process_job, connection={"host": "localhost", "port": 6379})

3. 同步代码阻塞异步事件循环

generate_report_task中调用的fetch_responses、responses_json_to_df等若为同步阻塞函数,会占用异步事件循环线程,导致Worker无法处理任务。

修复方法:用asyncio.to_thread将同步函数包装为异步执行:

async def generate_report_task(survey_id):
    try:
        logger.info(f"Starting report generation for survey ID: {survey_id}")
        # 将同步调用移到线程中执行,避免阻塞事件循环
        reponses, questions = await asyncio.to_thread(fetch_responses, survey_id)
        data_df = await asyncio.to_thread(responses_json_to_df, reponses, "definition")
        survey_data = await asyncio.to_thread(bdv_json_to_df, questions["definition"].iloc[0])

        output_path = f"output_{survey_id}.docx"
        # 同理处理同步生成报告的函数
        await asyncio.to_thread(
            decision_maker_with_report_streaming,
            data=data_df,
            df_bbl=survey_data,
            template="Template_.docx",
            bucket_name="reu-data",
            key=output_path
        )

        logger.info(f"Report generated successfully for survey ID: {survey_id}")
        return {"message": "Report generated successfully", "file_path": output_path}

    except Exception as e:
        logger.error(f"Error generating report for survey ID {survey_id}: {str(e)}")
        return {"error": str(e)}

4. 队列任务状态异常

任务可能被标记为延迟、暂停或失败,导致Worker无法拾取。可通过代码检查队列任务状态:

# 在Flask中添加检查路由
@app.route('/check-jobs', methods=['GET'])
async def check_jobs():
    waiting_jobs = await queue.getJobs(["waiting"])
    active_jobs = await queue.getJobs(["active"])
    failed_jobs = await queue.getJobs(["failed"])
    return jsonify({
        "waiting": len(waiting_jobs),
        "active": len(active_jobs),
        "failed": len(failed_jobs)
    })

若存在等待任务但Worker不处理,说明Worker与队列的连接或配置存在问题。

5. 线程异步循环兼容性问题

在单独线程中运行asyncio循环可能存在兼容性问题(尤其是Windows环境),可尝试将Worker与Flask分开启动:

  • 编写独立的worker.py脚本运行Worker,再单独启动Flask应用,避免同一进程内用线程隔离异步循环。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 21:47:03