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

