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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 23:45:10