如何使用MongoDB与Python实现股票数据实时更新时的均线交叉告警
Hey there! I get it—handling real-time stock data updates and triggering alerts on moving average crossovers can feel tricky when you're dealing with frequent inserts. Let's walk through practical solutions that fit your workflow, building on the Pandas work you've already done.
Core Approach Overview
Since you're inserting data every second, we need to avoid reprocessing the entire dataset every time. The key is to either:
- Use MongoDB's aggregation capabilities to fetch a sliding window of recent data and compute averages on the fly, or
- Maintain a local in-memory sliding window in Python to compute averages incrementally (faster for high-frequency inserts).
Option 1: MongoDB Aggregation for Real-Time Calculation
This is great if you want to offload some computation to the database and keep your Python logic clean. We'll use $setWindowFields (available in MongoDB 5.0+) to calculate rolling averages directly in an aggregation pipeline, then check for crossovers.
Example Code
from pymongo import MongoClient import datetime # Initialize MongoDB connection client = MongoClient("mongodb://localhost:27017/") db = client["stock_data"] collection = db["ohlc"] def insert_new_ohlc(ohlc_data): # Insert new data point (assuming ohlc_data has 'timestamp', 'close', etc.) collection.insert_one(ohlc_data) check_crossover() def check_crossover(): # Aggregation pipeline to get last 50 data points and compute moving averages pipeline = [ # Sort by timestamp to ensure we get the most recent data {"$sort": {"timestamp": -1}}, # Limit to last 50 points (adjust if your "day" uses different granularity) {"$limit": 50}, # Compute rolling 20 and 50-period averages using window functions {"$setWindowFields": { "sortBy": {"timestamp": 1}, "output": { "20daysMA": {"$avg": "$close", "window": {"documents": [-19, 0]}}, "50daysMA": {"$avg": "$close", "window": {"documents": [-49, 0]}} } }}, # Get the latest two data points to check for crossover {"$sort": {"timestamp": -1}}, {"$limit": 2} ] results = list(collection.aggregate(pipeline)) if len(results) < 2: return # Not enough data to check crossover latest = results[0] previous = results[1] # Check for golden cross (20MA crosses above 50MA) golden_cross = (previous["20daysMA"] <= previous["50daysMA"]) and (latest["20daysMA"] > latest["50daysMA"]) # Check for death cross (20MA crosses below 50MA) death_cross = (previous["20daysMA"] >= previous["50daysMA"]) and (latest["20daysMA"] < latest["50daysMA"]) if golden_cross: send_alert("Golden Cross Alert!", f"20-day MA crossed above 50-day MA at {latest['timestamp']}") elif death_cross: send_alert("Death Cross Alert!", f"20-day MA crossed below 50-day MA at {latest['timestamp']}") def send_alert(title, message): # Replace with your actual alert logic (email, Slack, etc.) print(f"ALERT: {title}\n{message}\n") # Example usage: Insert a new OHLC data point new_ohlc = { "timestamp": datetime.datetime.now(), "open": 150.2, "high": 151.5, "low": 149.8, "close": 150.7 } insert_new_ohlc(new_ohlc)
Option 2: Local In-Memory Sliding Window (Faster for High Frequency)
If you're dealing with very frequent inserts (like every second), maintaining a local sliding window with collections.deque reduces database round-trips. We'll keep the last 50 data points in memory, compute MAs with Pandas, and check crossovers.
Example Code
from pymongo import MongoClient import datetime from collections import deque import pandas as pd client = MongoClient("mongodb://localhost:27017/") db = client["stock_data"] collection = db["ohlc"] # Initialize sliding window to keep last 50 close prices window_size = 50 close_prices = deque(maxlen=window_size) # Track previous MA values to detect crossovers prev_20ma = None prev_50ma = None def insert_new_ohlc(ohlc_data): collection.insert_one(ohlc_data) # Update local window with new close price close_prices.append(ohlc_data["close"]) check_crossover() def check_crossover(): global prev_20ma, prev_50ma if len(close_prices) < 50: return # Need at least 50 points for 50-day MA # Convert deque to Pandas Series for MA calculation closes_series = pd.Series(close_prices) current_20ma = closes_series.rolling(window=20).mean().iloc[-1] current_50ma = closes_series.rolling(window=50).mean().iloc[-1] if prev_20ma is not None and prev_50ma is not None: golden_cross = (prev_20ma <= prev_50ma) and (current_20ma > current_50ma) death_cross = (prev_20ma >= prev_50ma) and (current_20ma < current_50ma) if golden_cross: send_alert("Golden Cross!", f"20MA ({current_20ma:.2f}) crossed above 50MA ({current_50ma:.2f})") elif death_cross: send_alert("Death Cross!", f"20MA ({current_20ma:.2f}) crossed below 50MA ({current_50ma:.2f})") # Update previous MA values for next check prev_20ma = current_20ma prev_50ma = current_50ma def send_alert(title, message): # Replace with your alert mechanism print(f"[{datetime.datetime.now()}] {title}\n{message}\n") # Example insert new_ohlc = { "timestamp": datetime.datetime.now(), "open": 150.2, "high": 151.5, "low": 149.8, "close": 150.7 } insert_new_ohlc(new_ohlc)
Key Notes to Avoid Issues
- Granularity Check: Make sure your "20-day" and "50-day" definitions match your data frequency. If you're inserting data every second, a "day" might mean 86400 points—adjust the window sizes accordingly!
- Deduplicate Alerts: Add a check to avoid sending multiple alerts for the same crossover (e.g., store the last crossover timestamp in MongoDB or a local variable).
- Error Handling: Add try/except blocks around database operations and alert logic to handle connection drops or data gaps.
内容的提问来源于stack exchange,提问作者Prakhar Gupta

