如何用Pandas处理MongoDB实时数据可视化及Flask联动统计更新
Hey there! Let's work through your real-time data visualization and stats update problem with Flask and MongoDB. First, let's fix a few small bugs in your existing Flask code to make sure it runs properly, then dive into solutions that will get your stats updating instantly when new data is added via Postman.
First: Fix Your Flask API Code
There are a couple of typos and configuration issues that will cause errors. Here's the corrected version:
from flask import Flask, jsonify, request from flask_pymongo import PyMongo app = Flask(__name__) # MongoDB default port is 27017 (you had 8000, which is likely your Flask port) app.config['MONGO_DBNAME'] = 'db' app.config['MONGO_URI'] = 'mongodb://localhost:27017/db' mongo = PyMongo(app) @app.route('/stocks', methods=['GET']) def get_all_stocks(): stocks = mongo.db.stocks output = [] for i in stocks.find(): output.append({'name': i['name'], 'item': i['item']}) return jsonify({'data': output}) @app.route('/add', methods=['POST']) def add_stocks(): stocks = mongo.db.stocks name = request.json.get('name') item = request.json.get('item') # Add validation for missing fields if not name or not item: return jsonify({'error': 'Missing required fields: name and item'}), 400 # Use insert_one (modern PyMongo syntax) and capture the correct ID item_id = stocks.insert_one({'name': name, 'item': item}).inserted_id new_stock = stocks.find_one({'_id': item_id}) output = {'name': new_stock['name'], 'item': new_stock['item']} return jsonify({'added': output}) # Fix route parameter syntax and typo in variable name @app.route('/stocks/<name>', methods=['GET']) def get_one_stock(name): stocks = mongo.db.stocks stock = stocks.find_one({'name': name}) if stock: output = {'name': stock['name'], 'item': stock['item']} else: output = 'No stock found with that name' return jsonify({'data': output}) if __name__ == '__main__': app.run(debug=True, port=8000)
Solutions for Real-Time Stats Updates
Running your Flask app and analysis script via a manager didn't work because they're separate processes that don't communicate. Here are three reliable approaches to sync them:
1. MongoDB Change Streams (Best Native Option)
MongoDB has built-in change streams that let you listen for real-time updates to your collection. This is the cleanest solution because it doesn't require extra tools.
Create your data_analysis_script.py like this:
from pymongo import MongoClient import pandas as pd import matplotlib.pyplot as plt from matplotlib.animation import FuncAnimation import threading # Connect to MongoDB client = MongoClient('mongodb://localhost:27017/') db = client['db'] stocks_collection = db['stocks'] # Get initial stats (e.g., count of items per name) def get_stock_stats(): pipeline = [ {'$group': {'_id': '$name', 'item_count': {'$sum': 1}}} ] stats = list(stocks_collection.aggregate(pipeline)) return pd.DataFrame(stats).rename(columns={'_id': 'name'}) # Set up visualization fig, ax = plt.subplots() ax.set_title('Stock Item Count by Name') ax.set_xlabel('Name') ax.set_ylabel('Number of Items') def update_chart(_): df = get_stock_stats() ax.clear() ax.bar(df['name'], df['item_count']) return ax, # Listen for MongoDB changes to trigger updates def watch_db_changes(): with stocks_collection.watch() as stream: for change in stream: # Refresh chart whenever data is inserted/updated plt.draw() plt.pause(0.01) if __name__ == '__main__': # Initialize chart with existing data initial_df = get_stock_stats() ax.bar(initial_df['name'], initial_df['item_count']) # Run change stream listener in a background thread watch_thread = threading.Thread(target=watch_db_changes, daemon=True) watch_thread.start() # Start auto-refreshing animation (fallback in case stream misses anything) ani = FuncAnimation(fig, update_chart, interval=1000) plt.show()
2. Flask Signals (Same Process Sync)
If you want to run your analysis script in the same process as Flask, use Flask signals to trigger updates when new data is added.
First, update your Flask app to include signals:
from flask import Flask, jsonify, request from flask_pymongo import PyMongo from blinker import Namespace # Create a signal for data updates signal_space = Namespace() data_added_signal = signal_space.signal('data-added') app = Flask(__name__) app.config['MONGO_DBNAME'] = 'db' app.config['MONGO_URI'] = 'mongodb://localhost:27017/db' mongo = PyMongo(app) # ... keep your other routes the same ... @app.route('/add', methods=['POST']) def add_stocks(): stocks = mongo.db.stocks name = request.json.get('name') item = request.json.get('item') if not name or not item: return jsonify({'error': 'Missing required fields'}), 400 item_id = stocks.insert_one({'name': name, 'item': item}).inserted_id new_stock = stocks.find_one({'_id': item_id}) # Send signal when data is added data_added_signal.send(app, data=new_stock) return jsonify({'added': {'name': new_stock['name'], 'item': new_stock['item']}})
Then your data_analysis_script.py can subscribe to the signal:
from your_flask_app import app, data_added_signal, mongo import pandas as pd import matplotlib.pyplot as plt import threading fig, ax = plt.subplots() def get_stock_stats(): stocks = mongo.db.stocks pipeline = [{'$group': {'_id': '$name', 'item_count': {'$sum': 1}}}] stats = list(stocks.aggregate(pipeline)) return pd.DataFrame(stats).rename(columns={'_id': 'name'}) def update_stats(sender, data): print(f"New data received: {data}") df = get_stock_stats() ax.clear() ax.bar(df['name'], df['item_count']) ax.set_title('Stock Item Count by Name') plt.draw() # Subscribe to the data added signal data_added_signal.connect(update_stats) if __name__ == '__main__': # Initialize chart initial_df = get_stock_stats() ax.bar(initial_df['name'], initial_df['item_count']) # Run Flask in a background thread flask_thread = threading.Thread(target=app.run, kwargs={'debug': True, 'port': 8000}, daemon=True) flask_thread.start() plt.show()
3. Redis Pub/Sub (Distributed Processes)
If you need to run Flask and your analysis script on separate servers or processes, use Redis as a message broker to send update triggers.
First, install Redis and the Python client: pip install redis
Update your Flask app to publish messages:
from flask import Flask, jsonify, request from flask_pymongo import PyMongo import redis # Connect to Redis r = redis.Redis(host='localhost', port=6379, db=0) app = Flask(__name__) app.config['MONGO_DBNAME'] = 'db' app.config['MONGO_URI'] = 'mongodb://localhost:27017/db' mongo = PyMongo(app) # ... keep your other routes the same ... @app.route('/add', methods=['POST']) def add_stocks(): stocks = mongo.db.stocks name = request.json.get('name') item = request.json.get('item') if not name or not item: return jsonify({'error': 'Missing required fields'}), 400 item_id = stocks.insert_one({'name': name, 'item': item}).inserted_id new_stock = stocks.find_one({'_id': item_id}) # Publish update message to Redis r.publish('stock_updates', 'new_data') return jsonify({'added': {'name': new_stock['name'], 'item': new_stock['item']}})
Then your data_analysis_script.py subscribes to Redis:
from pymongo import MongoClient import redis import pandas as pd import matplotlib.pyplot as plt import threading client = MongoClient('mongodb://localhost:27017/') db = client['db'] stocks_collection = db['stocks'] # Connect to Redis and subscribe to updates r = redis.Redis(host='localhost', port=6379, db=0) pubsub = r.pubsub() pubsub.subscribe('stock_updates') fig, ax = plt.subplots() def get_stock_stats(): pipeline = [{'$group': {'_id': '$name', 'item_count': {'$sum': 1}}}] stats = list(stocks_collection.aggregate(pipeline)) return pd.DataFrame(stats).rename(columns={'_id': 'name'}) def update_chart(): df = get_stock_stats() ax.clear() ax.bar(df['name'], df['item_count']) ax.set_title('Stock Item Count by Name') plt.draw() def listen_for_updates(): for message in pubsub.listen(): if message['type'] == 'message': update_chart() if __name__ == '__main__': initial_df = get_stock_stats() ax.bar(initial_df['name'], initial_df['item_count']) # Run Redis listener in background listen_thread = threading.Thread(target=listen_for_updates, daemon=True) listen_thread.start() plt.show()
Recommendation
I'd go with MongoDB Change Streams first—it's native to your database, requires no extra tools, and gives you precise control over what changes trigger updates.
内容的提问来源于stack exchange,提问作者Salman Al Farisi

