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

Celery Beat定时任务未按指定间隔触发问题求助

Celery Beat 未按指定间隔触发任务问题排查与解决

问题描述

开发Flask应用时,配置Celery Beat每10秒执行一次任务,但任务始终未触发。执行了以下命令:

  • python3 run.py
  • celery -A celery worker --loglevel=info
  • celery -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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 05:14:57