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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:58:17