如何获取Flask应用中Celery创建的未完成Redis任务列表
问题场景
我有两个运行Python Flask应用的容器(容器1和容器2):
- 容器1负责提交Celery任务,对应的Redis操作记录如下:
1677744862.693857 [0 192.168.80.13:44076] "SUBSCRIBE" "celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb" 1677744862.694585 [0 192.168.80.13:44044] "PING" 1677744862.694949 [0 192.168.80.13:44044] "LPUSH" "celery" "{\"body\": \"siDxNiIICwDfDwzDwiINcyWt5ZZd0WiCiwIEXg30Is2sGhtIons\", \"content-encoding\": \"utf-8\", \"content-type\": \"application/json\", \"headers\": {\"lang\": \"py\", \"task\": \"task.send_text\", \"id\": \"11902e64-94cc-4e02-89ce-3cea365294bb\", \"shadow\": null, \"eta\": null, \"expires\": null, \"group\": null, \"group_index\": null, \"retries\": 0, \"timelimit\": [null, null], \"root_id\": \"11902e64-94cc-4e02-89ce-3cea365294bb\", \"parent_id\": null, \"argsrepr\": \"('123456', '000_000_000_000', '000000', 'DefaultId', '', '1')\", \"kwargsrepr\": \"{}\", \"origin\": \"gen74@f0eb8cc30754\", \"ignore_result\": false}, \"properties\": {\"correlation_id\": \"11902e64-94cc-4e02-89ce-3cea365294bb\", \"reply_to\": \"7fdc8466-0de3-3be1-8a0a-a48713e2769a\", \"delivery_mode\": 2, \"delivery_info\": {\"exchange\": \"\", \"routing_key\": \"celery\"}, \"priority\": 0, \"body_encoding\": \"base64\", \"delivery_tag\": \"b5cc708a-8f28-4fb1-8714-1de4433c8621\"}}" 1677744862.696078 [0 192.168.80.13:44076] "UNSUBSCRIBE" "celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb"
- 容器2负责执行任务并写入结果,对应的Redis操作记录如下:
1677744993.925599 [0 192.168.80.12:43034] "GET" "celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb" 1677744993.927298 [0 192.168.80.12:43034] "MULTI" 1677744993.927340 [0 192.168.80.12:43034] "SETEX" "celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb" "86400" "{\"status\": \"SUCCESS\", \"result\": [\"True\", \"\", \"123456\"], \"traceback\": null, \"children\": [], \"date_done\": \"2023-03-02T08:16:33.922182\", \"task_id\": \"11902e64-94cc-4e02-89ce-3cea365294bb\",}" 1677744993.927480 [0 192.168.80.12:43034] "PUBLISH" "celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb" "{\"status\": \"SUCCESS\", \"result\": [\"True\", \"\", \"123456\"], \"traceback\": null, \"children\": [], \"date_done\": \"2023-03-02T08:16:33.922182\", \"task_id\": \"11902e64-94cc-4e02-89ce-3cea365294bb\",}" 1677744993.927602 [0 192.168.80.12:43034] "EXEC" 1677744994.868230 [0 192.168.80.12:42992] "BRPOP" "celery" "celery\x06\x163" "celery\x06\x166" "celery\x06\x169" "1"
现在容器2宕机,上述从GET开始的操作未执行,需要找出所有未完成的任务(比如celery-task-meta-11902e64-94cc-4e02-89ce-3cea365294bb)。但查看redis-cli info的# Keyspace部分,容器1提交的任务并未增加键数量,尝试以下Python代码也无法识别这类任务:
import redis import json redis_client = redis.Redis(host='localhost', port=6379, db=0) task_ids = redis_client.keys('celery-task-meta-*') unfinished_tasks = [] number_of_finished = [] for task_id in task_ids: task_result = redis_client.get(task_id) string = task_result.decode('utf-8') dictionary = json.loads(string) if "status" not in dictionary: unfinished_tasks.append(task_id) print('Unfinished tasks:', unfinished_tasks) print(len(unfinished_tasks))
问题原因
你的代码找不到未完成任务的核心原因是:celery-task-meta-<task_id>这个键只有在任务执行完成后才会被创建(由容器2的SETEX命令生成)。当容器2宕机未执行任务时,Redis中根本不存在这类键,自然无法通过keys('celery-task-meta-*')获取到。
未完成的任务实际上存储在Celery的任务队列中(默认是名为celery的列表),需要从队列中解析任务ID。
解决方案
方法1:直接解析Redis队列中的任务
从celery列表中取出所有未被消费的任务,解析每个任务的JSON内容,提取任务ID:
import redis import json redis_client = redis.Redis(host='localhost', port=6379, db=0) # 获取队列中所有未被消费的任务(LRANGE为只读操作,不会移除任务) queue_tasks = redis_client.lrange('celery', 0, -1) unfinished_task_ids = [] for task_bytes in queue_tasks: try: task_data = json.loads(task_bytes.decode('utf-8')) # 从任务headers中提取任务ID task_id = task_data['headers']['id'] unfinished_task_ids.append(task_id) except (json.JSONDecodeError, KeyError) as e: print(f"解析任务失败: {e}") continue print('未完成的任务ID:', unfinished_task_ids) print('未完成任务数量:', len(unfinished_task_ids))
方法2:使用Celery内置API查询
如果可以在容器1中直接调用Celery的API,更推荐这种方式,避免手动解析队列:
from celery import Celery # 初始化Celery实例,配置需与你的应用一致 app = Celery('tasks', broker='redis://localhost:6379/0') # 创建Celery检查实例,用于查询任务状态 inspect = app.control.inspect() # 获取队列中待执行的任务 reserved_tasks = inspect.reserved() # 获取正在执行的任务 active_tasks = inspect.active() unfinished_task_ids = [] # 解析待执行任务 if reserved_tasks: for worker, tasks in reserved_tasks.items(): for task in tasks: unfinished_task_ids.append(task['id']) # 解析正在执行的任务 if active_tasks: for worker, tasks in active_tasks.items(): for task in tasks: unfinished_task_ids.append(task['id']) print('未完成的任务ID:', unfinished_task_ids) print('未完成任务数量:', len(unfinished_task_ids))
补充说明
- 如果使用了多个任务队列(比如
celery\x06\x163这类带优先级的队列),需要遍历所有队列名称。可以通过redis_client.keys('celery*')获取所有队列键,再逐个解析。 - 方法1中的
LRANGE是只读操作,不会影响任务的正常消费;如果需要移除任务,可以使用LPOP或BRPOP,但需谨慎操作。
内容的提问来源于stack exchange,提问作者John
相关产品推荐
相关产品推荐

