Celery Beat定时任务未按指定间隔触发问题求助
Celery Beat 未按指定间隔触发任务问题排查与解决
问题描述
开发Flask应用时,配置Celery Beat每10秒执行一次任务,但任务始终未触发。执行了以下命令:
python3 run.pycelery -A celery worker --loglevel=infocelery -A celery worker --beat --loglevel=info
执行后出现的警告信息:
[2024-07-10 16:55:53,498: INFO/MainProcess] Connected to amqp://guest:**@127.0.0.1:6379/0// [2024-07-10 16:55:53,499: WARNING/MainProcess] /home/bhupesh/.local/lib/python3.10/site-packages/celery/worker/consumer/consumer.py:507: CPendingDeprecationWarning: The broker_connection_retry configuration setting will no longer determine whether broker connection retries are made during startup in Celery 6.0 and above. If you wish to retain the existing behavior for retrying connections on startup, you should set broker_connection_retry_on_startup to True. warnings.warn( [2024-07-10 16:55:53,534: INFO/MainProcess] mingle: searching for neighbors [2024-07-10 16:55:53,543: WARNING/MainProcess] No hostname was supplied. Reverting to default 'localhost' [2024-07-10 16:55:54,571: INFO/MainProcess] mingle: all alone [2024-07-10 16:55:54,624: INFO/MainProcess] celery@vijay-latitude-7480 ready.
相关代码:
from flask import Flask, request, jsonify from flask_sqlalchemy import SQLAlchemy from datetime import datetime import requests import xmltodict import json import logging import os from celery.schedules import crontab from celery import Celery from .celery import make_celery from datetime import timedelta db = SQLAlchemy() class Event(db.Model): id = db.Column(db.Integer, primary_key=True) provider_id = db.Column(db.Integer, unique=False, nullable=False) name = db.Column(db.String(255), nullable=False) started_at = db.Column(db.DateTime, nullable=False) end_at = db.Column(db.DateTime, nullable=False) sell_mode = db.Column(db.String(50), nullable=False) created_at = db.Column(db.DateTime, nullable=False, default=db.func.current_timestamp()) def as_dict(self): """ Converts the Event object to a dictionary representation. Returns: dict: Dictionary representation of the Event object. """ return { "provider_id": self.provider_id, "name": self.name, "datetime": self.started_at.isoformat(), "sell_mode": self.sell_mode } def remove_at_sign(obj): """ Recursively removes '@' sign from keys """ if isinstance(obj, dict): return {key.lstrip('@'): remove_at_sign(val) for key, val in obj.items()} elif isinstance(obj, list): return [remove_at_sign(elem) for elem in obj] else: return obj def fetch_events(): """ Fetches events data from an external API and processes it. Returns: list: List of Event objects fetched from the API. """ PROVIDER_URL = 'some-url' response = requests.get(PROVIDER_URL) events = [] if response.status_code == 200: json_resp = xmltodict.parse(response.text) final_json = remove_at_sign(json_resp) with open('events.json', 'a') as f1: json.dump(final_json, f1, indent=4) final_data = final_json['eventList']['output']['base_event'] for ind, data in enumerate(final_data): if 'sell_mode' in data and final_data[ind]['sell_mode'] == 'online': sell_mode = final_data[ind]['sell_mode'] title = data['title'] if 'title' in data else str() start_date = datetime.fromisoformat(data['event']['event_start_date']) end_date = datetime.fromisoformat(data['event']['event_end_date']) event_id = data['event']['event_id'] events.append(Event(provider_id=event_id, name=title, end_at=end_date, sell_mode=sell_mode, created_at=datetime.now(), started_at=start_date)) print([sell_mode, title, start_date, end_date, event_id]) return events def store_events(events, app): """ Stores events in the database. Args: events (list): List of Event objects to be stored. app (Flask): Flask application instance. """ with app.app_context(): try: db.session.add_all(events) db.session.commit() logging.info(f"Stored {len(events)} new events in the database.") except Exception as e: db.session.rollback() logging.error(f"Failed to store events: {str(e)}") raise def fetch_and_store(app): with app.app_context(): db.create_all() events = fetch_events() store_events(events, app) def create_app(): # import ipdb;ipdb.set_trace() app = Flask(__name__) app.config['SQLALCHEMY_DATABASE_URI'] = 'sqlite:///events.db' app.config['SQLALCHEMY_TRACK_MODIFICATIONS'] = False app.config.update( CELERY_BROKER_URL='redis://127.0.0.1:6379/0', CELERY_RESULT_BACKEND='redis://127.0.0.1:6379/0' ) db.init_app(app) celery = make_celery(app) celery_worker = 'worker1@hostname' # Configure logging logging.basicConfig(level=logging.DEBUG, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s') log_file = os.path.join(app.root_path, 'app.log') file_handler = logging.FileHandler(log_file) file_handler.setLevel(logging.INFO) file_handler.setFormatter(logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')) logging.getLogger().addHandler(file_handler) logger = logging.getLogger(__name__) with app.app_context(): db.create_all() @celery.task(name='app.scheduled_fetch_and_store') def scheduled_fetch_and_store(): fetch_and_store(app) celery.conf.beat_schedule = { 'fetch-events-every-10-seconds': { 'task': 'app.scheduled_fetch_and_store', 'schedule': timedelta(seconds=10), }, } @app.route('/events', methods=['GET']) def get_events(): """ Retrieves events filtered by start and end dates. Returns: jsonify: JSON response containing events data. """ starts_at = request.args.get('starts_at') ends_at = request.args.get('ends_at') logger.info(f"Fetching events between {starts_at} and {ends_at}") if not starts_at or not ends_at: logger.error("Both 'starts_at' and 'ends_at' parameters are required") return jsonify({"error": "Both 'starts_at' and 'ends_at' parameters are required"}), 400 try: starts_at = datetime.fromisoformat(starts_at) ends_at = datetime.fromisoformat(ends_at) except ValueError: logger.error("'starts_at' and 'ends_at' must be in ISO format") return jsonify({"error": "'starts_at' and 'ends_at' must be in ISO format"}), 400 events = Event.query.filter(Event.started_at >= starts_at, Event.end_at <= ends_at).all() result = [ { "id": event.provider_id, "name": event.name, "start_time": event.started_at.isoformat(), "end_time": event.end_at.isoformat(), "sell_mode": event.sell_mode } for event in events ] with open('sample.json', 'w') as json_file: json.dump(result, json_file, indent=4) return jsonify(result) @app.route('/routes', methods=['GET']) def list_routes(): import urllib.parse output = [] for rule in app.url_map.iter_rules(): options = {} for arg in rule.arguments: options[arg] = f"[{arg}]" methods = ','.join(rule.methods) url = urllib.parse.unquote(f"{rule}") line = f"{rule.endpoint:50s} {methods:20s} {url}" output.append(line) return "<pre>" + "\n".join(sorted(output)) + "</pre>" return app if __name__ == "__main__": app = create_app() app.run(debug=True)
排查与解决方案
1. 修正Celery启动命令
你当前的启动命令celery -A celery worker --beat没有指向正确的Flask应用模块,导致Beat无法加载create_app内的调度配置。
假设你的入口文件是run.py,修改启动命令为:
celery -A run.celery worker --beat --loglevel=info
前提是在run.py中导出了celery实例(可参考下文代码调整)。
2. 调整Celery配置加载时机
当前Beat调度配置放在create_app内部的app.app_context()块中,只有Flask app启动时才会加载,Celery单独启动时无法读取到配置。
修改代码,将Celery实例和配置导出到全局:
# 在run.py顶部添加全局celery实例声明 celery = None def create_app(): global celery app = Flask(__name__) # ... 原有配置代码 ... celery = make_celery(app) # 移到app_context外部定义任务和调度 @celery.task(name='app.scheduled_fetch_and_store') def scheduled_fetch_and_store(): fetch_and_store(app) celery.conf.beat_schedule = { 'fetch-events-every-10-seconds': { 'task': 'app.scheduled_fetch_and_store', 'schedule': timedelta(seconds=10), }, } # ... 原有代码 ... return app
3. 修复Broker连接错误
日志显示Celery尝试连接AMQP协议(默认RabbitMQ),但你配置的是Redis作为Broker,说明make_celery函数未正确读取Flask配置。
修正make_celery实现:
def make_celery(app): celery = Celery( app.name, broker=app.config['CELERY_BROKER_URL'], backend=app.config['CELERY_RESULT_BACKEND'] ) celery.conf.update(app.config) # 添加Flask上下文支持 class ContextTask(celery.Task): def __call__(self, *args, **kwargs): with app.app_context(): return self.run(*args, **kwargs) celery.Task = ContextTask return celery
4. 解决警告信息
- 针对
broker_connection_retry警告:在Flask配置中添加CELERY_BROKER_CONNECTION_RETRY_ON_STARTUP = True - 针对
No hostname was supplied警告:启动时指定hostname参数:celery -A run.celery worker --beat --loglevel=info --hostname=worker1@localhost
内容的提问来源于stack exchange,提问作者Vijay Bokade
相关产品推荐
相关产品推荐

