如何用MongoDB聚合OHLC K线?生成多周期K线及数据过滤
Got it, let's break down how to build a function that converts your 1-minute OHLC data into any desired timeframe (5min, 1h, daily, etc.) while supporting data filtering. I'll use MongoDB's aggregation framework since it's perfect for this kind of data rollup, and wrap it into a reusable function (I'll use Python with pymongo as an example, but the logic translates to other drivers too).
Core Approach
The key steps are straightforward:
- Combine Date and Time into a single datetime field: MongoDB's date functions work best with a unified ISODate, so we'll merge your separate
DateandTimestrings into one timestamp. - Apply filters (optional): Narrow down the dataset first to speed up processing—think date ranges, volume thresholds, etc.
- Group by target timeframe: Bucket 1-minute bars into your desired period, then calculate OHLCV values for each bucket.
- Sort results: Ensure the output K-lines are in chronological order for usability.
Step 1: Aggregation Pipeline Breakdown
Let's walk through each stage of the pipeline that does the heavy lifting:
1.1 Filter Stage (Optional)
Use $match to trim down your data before rolling it up. For example, filter to a specific date range or exclude zero-volume bars:
{ $match: { "Date": { "$gte": "08/06/2007", "$lte": "08/07/2007" }, "Volume": { "$gt": 0 } } }
1.2 Merge Date and Time into Datetime
We'll convert your separate date/time strings into a single ISODate using $dateFromString. Add this as an $addFields stage:
{ $addFields: { "datetime": { $dateFromString: { dateString: { $concat: ["$Date", " ", "$Time"] }, format: "%m/%d/%Y %H:%M", timezone: "UTC" // Adjust to your data's timezone if needed } } } }
1.3 Group by Target Timeframe
Use $dateTrunc (MongoDB 5.0+) to bucket the datetime into your desired period. For older MongoDB versions, we'll use a fallback with timestamp math.
For example, grouping into 5-minute buckets:
{ $group: { _id: { $dateTrunc: { date: "$datetime", unit: "minute", binSize: 5, timezone: "UTC" } }, Open: { $first: "$Open" }, // First tick in the bucket High: { $max: "$High" }, // Highest price in the bucket Low: { $min: "$Low" }, // Lowest price in the bucket Close: { $last: "$Close" }, // Last tick in the bucket Up: { $sum: "$Up" }, // Sum all Up values Down: { $sum: "$Down" }, // Sum all Down values Volume: { $sum: "$Volume" } // Sum all Volume values } }
For MongoDB versions <5.0, replace the _id with this timestamp math (example for 5 minutes):
_id: { $toDate: { $multiply: [ { $floor: { $divide: [{ $toLong: "$datetime" }, 300000] } }, // 5min = 300,000 ms 300000 ] } }
1.4 Sort the Results
Finally, sort by the grouped datetime to get a proper K-line sequence:
{ $sort: { "_id": 1 } }
Step 2: Reusable Function (Python + Pymongo)
Let's wrap this into a function that's easy to use with different timeframes and filters:
from pymongo import MongoClient def generate_custom_ohlc(collection, target_period, filter_query=None): """ Generate custom timeframe OHLC data from 1-minute bars in MongoDB. Args: collection: Pymongo collection object with your 1-minute OHLC data. target_period: String for the desired timeframe (e.g., "5min", "1h", "1d"). filter_query: Optional MongoDB match query to filter data before rollup. Returns: List of custom timeframe OHLC documents. """ # Map target periods to dateTrunc parameters and millisecond values period_map = { "1min": {"unit": "minute", "binSize": 1, "ms": 60000}, "5min": {"unit": "minute", "binSize": 5, "ms": 300000}, "15min": {"unit": "minute", "binSize": 15, "ms": 900000}, "30min": {"unit": "minute", "binSize": 30, "ms": 1800000}, "1h": {"unit": "hour", "binSize": 1, "ms": 3600000}, "4h": {"unit": "hour", "binSize": 4, "ms": 14400000}, "1d": {"unit": "day", "binSize": 1, "ms": 86400000} } if target_period not in period_map: raise ValueError(f"Unsupported period: {target_period}. Use one of {list(period_map.keys())}") period_params = period_map[target_period] pipeline = [] # Add filter stage if provided if filter_query: pipeline.append({"$match": filter_query}) # Merge Date and Time into a unified datetime field pipeline.append({ "$addFields": { "datetime": { "$dateFromString": { "dateString": {"$concat": ["$Date", " ", "$Time"]}, "format": "%m/%d/%Y %H:%M", "timezone": "UTC" } } } }) # Build group stage with fallback for older MongoDB versions try: group_id = { "$dateTrunc": { "date": "$datetime", "unit": period_params["unit"], "binSize": period_params["binSize"], "timezone": "UTC" } } except: group_id = { "$toDate": { "$multiply": [ {"$floor": {"$divide": [{"$toLong": "$datetime"}, period_params["ms"]]}}, period_params["ms"] ] } } pipeline.append({ "$group": { "_id": group_id, "Open": {"$first": "$Open"}, "High": {"$max": "$High"}, "Low": {"$min": "$Low"}, "Close": {"$last": "$Close"}, "Up": {"$sum": "$Up"}, "Down": {"$sum": "$Down"}, "Volume": {"$sum": "$Volume"} } }) # Sort results chronologically pipeline.append({"$sort": {"_id": 1}}) # Execute aggregation and return results return list(collection.aggregate(pipeline)) # Example usage client = MongoClient("mongodb://localhost:27017/") db = client["your_database_name"] ohlc_collection = db["your_1min_ohlc_collection"] # Generate 5-minute OHLC data filtered to 08/06/2007 five_min_ohlc = generate_custom_ohlc( ohlc_collection, "5min", filter_query={"Date": "08/06/2007"} ) # Print the first 2 bars to verify for bar in five_min_ohlc[:2]: print(bar)
Key Notes
- Timezone: Adjust the
timezoneparameter in$dateFromStringand$dateTruncto match your data's actual timezone (e.g., "America/New_York"). - Custom Periods: Extend the
period_mapwith your own timeframes (e.g., "2h" by adding{"unit": "hour", "binSize":2, "ms":7200000}). - Performance: For large datasets, add an index on the
datetimefield after creating it:collection.create_index("datetime")to speed up aggregation. - Gap Filling: If you need to include empty buckets (where no 1-minute data exists), you'll need an extra step to generate those empty periods and merge them with your results.
内容的提问来源于stack exchange,提问作者jersey bean

