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

如何同时运行Flask RESTful API服务器并存储Message Broker传入的JSON数据?

解决方案:同时运行Flask API与消息消费任务

这个问题其实很常见——核心就是要让Flask API服务和消息消费两个任务并行运行,同时保证数据操作的安全性。我给你一套完整的实现方案,一步步来:

1. 先导入需要的依赖

首先把要用的库都导入进来,这里以RabbitMQ作为Message Broker示例(你可以换成自己在用的Kafka、Redis等对应的客户端):

from flask import Flask, jsonify
import threading
import json
import pika  # 换成你的Broker客户端,比如kafka-python如果用Kafka

2. 初始化全局变量

定义存储数据的字典,还有线程锁(这个很重要,避免多线程同时读写字典出问题):

app = Flask(__name__)

# 用来存消息数据的全局字典
message_store = {}
# 线程锁,保证多线程下字典操作的安全性
store_lock = threading.Lock()

3. 写消息消费的逻辑

把消费Broker消息、解析JSON并存入字典的逻辑封装成一个函数,记得操作字典时加锁:

def consume_messages():
    # 连接你的Message Broker,这里是RabbitMQ的示例配置
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()
    
    # 声明你要消费的队列(根据实际情况调整名称)
    channel.queue_declare(queue='your_target_queue')
    
    def callback(ch, method, properties, body):
        try:
            # 把字节流解析成JSON
            message_data = json.loads(body.decode('utf-8'))
            # 假设消息里有唯一标识,比如'id',用它当字典的键
            item_id = message_data.get('id')
            if item_id:
                # 加锁后再操作字典,防止多线程冲突
                with store_lock:
                    message_store[item_id] = message_data
                print(f"已存储消息ID: {item_id}")
            else:
                print("消息缺少必填的'id'字段,跳过存储")
        except json.JSONDecodeError:
            print("收到无效的JSON格式消息")
    
    # 开始持续消费消息
    channel.basic_consume(queue='your_target_queue', on_message_callback=callback, auto_ack=True)
    print("消息消费线程已启动...")
    channel.start_consuming()

4. 实现Flask API查询端点

写两个示例端点:一个根据ID查单条消息,一个查所有消息,同样操作字典时要加锁:

@app.route('/message/<string:msg_id>', methods=['GET'])
def get_single_message(msg_id):
    # 加锁读取字典,避免读取过程中数据被修改
    with store_lock:
        target_msg = message_store.get(msg_id)
    if target_msg:
        return jsonify(target_msg), 200
    else:
        return jsonify({"error": "未找到该ID的消息"}), 404

# 可选:获取所有已存储的消息
@app.route('/messages', methods=['GET'])
def get_all_messages():
    with store_lock:
        # 返回字典的副本,防止外部直接修改原数据
        messages_copy = message_store.copy()
    return jsonify(messages_copy), 200

5. 启动两个任务

在主入口里,先启动消息消费的线程,再启动Flask服务器:

if __name__ == '__main__':
    # 创建消息消费线程
    consumer_thread = threading.Thread(target=consume_messages)
    # 设置为守护线程:当Flask主线程退出时,消费线程自动终止
    consumer_thread.daemon = True
    consumer_thread.start()
    
    # 启动Flask服务器(生产环境别用debug=True,这里只是示例)
    app.run(host='0.0.0.0', port=5000, debug=False)

几个关键注意事项

  • 线程安全不能忘:一定要用threading.Lock保护对message_store的读写,不然多线程同时操作字典很容易出现数据错乱或者崩溃。
  • 守护线程设置:把消费线程设为守护线程,这样关闭Flask服务器时,消费线程会跟着退出,不会留下后台僵尸进程。
  • Broker适配:如果用的不是RabbitMQ,只需要把consume_messages里的连接和消费逻辑换成对应Broker的客户端代码就行(比如Kafka的consumer.poll()循环)。
  • 生产环境部署:别用app.run()跑Flask,换成Gunicorn、uWSGI这类专业的WSGI服务器,还可以用Supervisor之类的进程管理器来监控两个任务的运行状态。

这样一来,Flask API会在前台处理HTTP查询请求,消息消费线程在后台持续接收消息并存入字典,两者并行运行,数据操作也安全可靠。

内容的提问来源于stack exchange,提问作者Yash Mishra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 12:28:12