使用Mongoose聚合MongoDB多文档:按时间框架合并为单文档
MongoDB分钟级K线数据按参数聚合实现方案
需求说明
- 集合存储数千条1分钟级K线文档,每条含唯一
timestamp字段 - 根据API传入的
timeframe和limit参数执行聚合:timeframe: '1m'+limit: 30:取最近30条1分钟数据,按需求聚合(可生成30条1m K线或单条汇总文档)timeframe: '2m'+limit: 30:取最近60条1分钟数据,聚合生成30条2m K线timeframe: '5m'+limit: 10:取最近50条1分钟数据,聚合生成10条5m K线
文档结构示例:
{ "_id" : ObjectId("63049944031441c5312f4954"), "name" : "kucoin:ADA/USDT", "exchange" : "kucoin", "timeframe" : "1m", "open" : 0.451433, "high" : 0.451777, "low" : 0.451433, "close" : 0.451585, "volume" : 1042.7538, "timestamp" : 1661185920000.0, "__v" : 0 }
实现方案
1. 计算需获取的原始数据条数
先根据参数组合确定要拉取的1分钟数据总量:
const getFetchCount = (timeframe, limit) => { const multiplier = { '1m':1, '2m':2, '5m':5 }[timeframe]; return limit * multiplier; }; // 示例:getFetchCount('2m',30) → 60
2. 聚合生成目标周期K线(多文档输出)
如果需要生成limit条目标周期的K线文档(比如30条2m K线),用以下聚合管道:
// 定义传入的参数 const timeframe = '2m'; const limit = 30; const interval = { '1m':60000, '2m':120000, '5m':300000 }[timeframe]; const fetchCount = getFetchCount(timeframe, limit); db.collection.aggregate([ // 取最近fetchCount条1分钟数据,按时间升序排列 { $sort: { timestamp: -1 } }, { $limit: fetchCount }, { $sort: { timestamp: 1 } }, // 计算每条数据所属的目标周期时间戳(向下取整到周期起点) { $addFields: { targetTs: { $subtract: ['$timestamp', { $mod: ['$timestamp', interval] }] } } }, // 按目标周期分组,计算K线核心指标 { $group: { _id: { targetTs: '$targetTs', name: '$name', exchange: '$exchange', timeframe: timeframe }, open: { $first: '$open' }, // 周期第一条的开盘价 high: { $max: '$high' }, // 周期内最高价 low: { $min: '$low' }, // 周期内最低价 close: { $last: '$close' }, // 周期最后一条的收盘价 volume: { $sum: '$volume' } // 周期内成交量总和 } }, // 整理输出字段,按时间降序取limit条 { $project: { _id: 0, name: '$_id.name', exchange: '$_id.exchange', timeframe: '$_id.timeframe', timestamp: '$_id.targetTs', open: 1, high: 1, low: 1, close: 1, volume: 1 } }, { $sort: { timestamp: -1 } }, { $limit: limit } ])
3. 聚合为单条汇总文档(单文档输出)
如果需求是将所有取到的原始数据合并为单条汇总文档(比如统计这段时间的整体行情),用以下管道:
const timeframe = '1m'; const limit = 30; const fetchCount = getFetchCount(timeframe, limit); db.collection.aggregate([ { $sort: { timestamp: -1 } }, { $limit: fetchCount }, { $sort: { timestamp: 1 } }, // 全部数据归为一组,计算汇总指标 { $group: { _id: null, name: { $first: '$name' }, exchange: { $first: '$exchange' }, timeframe: timeframe, startTime: { $first: '$timestamp' }, endTime: { $last: '$timestamp' }, open: { $first: '$open' }, high: { $max: '$high' }, low: { $min: '$low' }, close: { $last: '$close' }, totalVolume: { $sum: '$volume' }, dataCount: { $sum: 1 } } }, // 整理输出字段 { $project: { _id: 0, name: 1, exchange: 1, timeframe: 1, startTime: 1, endTime: 1, open: 1, high: 1, low: 1, close: 1, totalVolume: 1, dataCount: 1 } } ])
内容的提问来源于stack exchange,提问作者Ahmed Kabeer
相关产品推荐
相关产品推荐

