使用RabbitMQ、Celery与Flask实现数据库更新时Celery无法识别消息如何解决
问题根因
你遇到的报错核心是:直接使用pika原生客户端向Celery队列推送自定义JSON消息,不符合Celery的消息协议格式。Celery的任务消息包含固定的元信息(任务名、参数、任务ID等),纯自定义JSON会被Celery识别为未知消息直接丢弃。
其他存在的问题
- 不需要手动维护RabbitMQ连接和消息发送逻辑,直接调用Celery提供的任务API即可完成异步任务推送
- 代码中
INSERT_Query、updateQuery都是占位字符串,没有替换为实际可执行的SQL,执行时会报错 - 全局初始化的RabbitMQ长连接会存在超时断开问题,手动维护容易引发连接异常
- 配置没有统一复用,两边硬编码配置容易出现不一致
- sqlite3执行写操作后没有手动提交事务,修改不会持久化到数据库
- 建表语句没有加
IF NOT EXISTS判断,重复调用接口会报错
修改后参考代码
consumer.py
from celery import Celery import sqlite3 import time import configparser # 加载统一配置 parser = configparser.RawConfigParser() parser.read('appconfig.conf') rmq_username = parser.get('general', 'rmq_USERNAME') rmq_password = parser.get('general', 'rmq_PASSWORD') host = parser.get('general', 'rmq_IP') port = parser.get('general', 'rmq_PORT') DATABASE = parser.get('general', 'DATABASE_FILE') broker_url = f'pyamqp://{rmq_username}:{rmq_password}@{host}:{port}//' app = Celery('tasks', backend='rpc://', broker=broker_url) @app.task(serializer='json') def updateDB(x): x = x["item"] with sqlite3.connect(DATABASE) as conn: time.sleep(5) # 可根据实际业务调整更新逻辑 conn.execute('''UPDATE table1 SET feild3='completed' WHERE feild2=?''', (x,)) conn.commit() return x
ProcedureAPI.py
from flask import Flask,request,jsonify import pandas as pd import sqlite3 import configparser # 导入Celery异步任务 from consumer import updateDB parser = configparser.RawConfigParser() configFilePath = 'appconfig.conf' parser.read(configFilePath) DATABASE= parser.get('general', 'DATABASE_FILE') app = Flask(__name__) @app.route('/create', methods=['POST']) def create_main(): if request.method=="POST": with sqlite3.connect(DATABASE) as conn: conn.execute('''CREATE TABLE IF NOT EXISTS table1 (feild1 INTEGER PRIMARY KEY AUTOINCREMENT, feild2 varchar(20) NOT NULL, feild3 varchar(20) DEFAULT 'pending');''') conn.commit() return "Table created",202 @app.route('/getData', methods=['GET']) def display_main(): if request.method=="GET": with sqlite3.connect(DATABASE) as conn: df = pd.read_sql_query("SELECT * from table1", conn) df_list = df.values.tolist() return jsonify(df_list),200 @app.route('/', methods=['POST']) def update_main(): if request.method=="POST": updatedata=request.get_json() with sqlite3.connect(DATABASE) as conn: # 可根据实际业务调整插入逻辑 conn.execute("INSERT INTO table1(feild2) VALUES (?)", (updatedata['item'],)) conn.commit() print("Records Inserted successfully") # 推送异步任务,Celery会自动封装消息发送到RabbitMQ updateDB.delay(updatedata) return updatedata,202 if __name__ == '__main__': app.run(debug=True)
启动&测试说明
- Celery启动命令无需修改,仍然使用:
celery -A consumer worker --pool=solo -l info - 测试时POST请求体传入
{"item": "实际值"}格式的JSON即可触发完整流程
内容的提问来源于stack exchange,提问作者Mahipal Singh
相关产品推荐
相关产品推荐

