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
相关产品推荐
相关产品推荐

