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

如何在Celery Beat任务参数中传入城市列表批量执行任务?

问题描述

我有一个包含10个城市的列表,想让Celery Beat每隔1分钟对每个城市执行get_weather任务,但当前代码无法实现预期效果,相关代码如下:

cityList = ["Tehran", "Shiraz", "Mashhad", "Qom", "Isfahan", "Ardabil", "Hamedan", "Yazd", "Tabriz", "Zavareh"]

app.conf.beat_schedule = {
    'call_show_every_one_minute': {
        "task": 'ali111.get_weather',
        'schedule': crontab(minute='*/1'),
        'args': ([cityList], ),
    }
}

@app.task()
def get_weather(city):
    con = redis.StrictRedis(host='localhost',port=6380, db=0, decode_responses=True)
    appid = "b3bf68fdfc6ba46923cd50cb8b9a79c3"
    URL = 'https://api.openweathermap.org/data/2.5/weather'
    temp = con.get(city)
    if temp is not None:
        return temp
    try:
        PARAMS = {'q' :city, 'appid' :appid}
        r = requests.get(url=URL, params=PARAMS)
        city_temp = (r.json()['main']['temp']) - 273.15
        my_temp = f"{round(city_temp,1)}c"
        con.set(city, my_temp, ex=60)
        return my_temp
    except ConnectionError:
        return "not internet connection"
问题原因

当前代码无法实现需求的核心问题:

  • beat_schedule中错误地将整个cityList作为单个参数传入get_weather任务,但该任务仅接受单个城市字符串作为参数,导致参数不匹配,任务无法正常执行。
  • 当前配置只会每分钟触发一次任务,而非为列表中的每个城市单独触发任务。
解决方法

方法一:为每个城市单独添加Beat调度任务

循环遍历城市列表,给每个城市创建独立的Beat任务配置,这样每分钟每个城市都会触发一次get_weather任务:

cityList = ["Tehran", "Shiraz", "Mashhad", "Qom", "Isfahan", "Ardabil", "Hamedan", "Yazd", "Tabriz", "Zavareh"]

# 初始化beat_schedule
app.conf.beat_schedule = {}

# 为每个城市添加调度任务
for city in cityList:
    app.conf.beat_schedule[f'get_weather_{city.lower()}'] = {
        "task": 'ali111.get_weather',
        'schedule': crontab(minute='*/1'),
        'args': (city, ),
    }

@app.task()
def get_weather(city):
    con = redis.StrictRedis(host='localhost',port=6380, db=0, decode_responses=True)
    appid = "b3bf68fdfc6ba46923cd50cb8b9a79c3"
    URL = 'https://api.openweathermap.org/data/2.5/weather'
    temp = con.get(city)
    if temp is not None:
        return temp
    try:
        PARAMS = {'q' :city, 'appid' :appid}
        r = requests.get(url=URL, params=PARAMS)
        city_temp = (r.json()['main']['temp']) - 273.15
        my_temp = f"{round(city_temp,1)}c"
        con.set(city, my_temp, ex=60)
        return my_temp
    except ConnectionError:
        return "not internet connection"

方法二:创建批量任务,一次性触发所有城市的天气查询

如果不想创建多个Beat任务,可以编写一个批量任务,Beat每分钟触发这个批量任务,由它来调用每个城市的get_weather任务:

from celery import group

cityList = ["Tehran", "Shiraz", "Mashhad", "Qom", "Isfahan", "Ardabil", "Hamedan", "Yazd", "Tabriz", "Zavareh"]

app.conf.beat_schedule = {
    'get_all_weather_every_one_minute': {
        "task": 'ali111.get_all_weather',
        'schedule': crontab(minute='*/1'),
    }
}

@app.task()
def get_weather(city):
    con = redis.StrictRedis(host='localhost',port=6380, db=0, decode_responses=True)
    appid = "b3bf68fdfc6ba46923cd50cb8b9a79c3"
    URL = 'https://api.openweathermap.org/data/2.5/weather'
    temp = con.get(city)
    if temp is not None:
        return temp
    try:
        PARAMS = {'q' :city, 'appid' :appid}
        r = requests.get(url=URL, params=PARAMS)
        city_temp = (r.json()['main']['temp']) - 273.15
        my_temp = f"{round(city_temp,1)}c"
        con.set(city, my_temp, ex=60)
        return my_temp
    except ConnectionError:
        return "not internet connection"

@app.task()
def get_all_weather():
    # 创建任务组,批量调用get_weather
    job = group(get_weather.s(city) for city in cityList)
    job.apply_async()
注意事项
  • 方法一的好处是每个城市的任务可以独立监控和重试,缺点是当城市列表变动时需要修改代码并重新启动Beat。
  • 方法二更灵活,城市列表变动时只需修改cityList即可,但所有城市任务会作为一个组执行,监控粒度为整个任务组。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 03:21:36