如何在InfluxDB中从Tick数据计算1分钟OHLCV数据?
从InfluxDB Tick数据计算1分钟OHLCV的Flux脚本修复
背景信息
我在InfluxDB中存储了期货Tick数据,行协议格式如下(lastPrice为实时价格,totalVolume为当日累计成交量):
tick,contract=FU2301 lastPrice=10.2,totalVolume=100 tick,contract=FU2301 lastPrice=10.1,totalVolume=110 tick,contract=FU2301 lastPrice=10.3,totalVolume=150 ... tick,contract=FU2301 lastPrice=9.8,totalVolume=290
需要基于这些Tick数据计算1分钟OHLCV数据,其中volume为该分钟的成交量(非累计值),目标格式如下:
time measurement contract open high low close volume ---- xxxx tick FU2301 10.2 10.3 9.0 9.3 10 xxxx tick FU2301 10.2 10.3 9.1 9.2 40 xxxx tick FU2301 10.1 10.4 9.0 9.1 20
我编写了Flux脚本,但volume计算部分未达到预期效果,原脚本如下:
data = from(bucket: "mk_data_test") |> range(start: -12h, stop: now()) |> filter(fn: (r) => r["_measurement"] == "tick") |> filter(fn: (r) => r["contract"] == "FU2301") |> window(every: 1m) volumeData = data |> filter(fn: (r) => r["_field"] == "totalVolume") priceData = data |> filter(fn: (r) => r["_field"] == "lastPrice") max = priceData |> max() |> set(key: "_field", value: "max") min = priceData |> min() |> set(key: "_field", value: "min") open = priceData |> first() |> set(key: "_field", value: "open") close = priceData |> last() |> set(key: "_field", value: "close") // 此处volume计算不符合预期 volume = volumeData |> first() |> difference() |> set(key: "_field", value: "volume") union(tables: [max, min, open, close, volume]) |> pivot(rowKey: ["_start"], columnKey: ["_field"], valueColumn: "_value")
问题分析
原脚本中volume计算逻辑错误:
- 先对每个1分钟窗口取
totalVolume的第一个值,再执行difference(),得到的是相邻窗口起始累计成交量的差值,而非当前分钟内的成交量。 - 正确逻辑:每个分钟窗口内,用窗口最后一条数据的
totalVolume减去窗口第一条数据的totalVolume,得到该分钟的实际成交量。
修复后的脚本
data = from(bucket: "mk_data_test") |> range(start: -12h, stop: now()) |> filter(fn: (r) => r["_measurement"] == "tick") |> filter(fn: (r) => r["contract"] == "FU2301") |> window(every: 1m) volumeData = data |> filter(fn: (r) => r["_field"] == "totalVolume") priceData = data |> filter(fn: (r) => r["_field"] == "lastPrice") max = priceData |> max() |> set(key: "_field", value: "max") min = priceData |> min() |> set(key: "_field", value: "min") open = priceData |> first() |> set(key: "_field", value: "open") close = priceData |> last() |> set(key: "_field", value: "close") // 修复volume计算:窗口内last totalVolume减去first totalVolume volume = volumeData |> aggregateWindow(every: 1m, fn: (tables) => tables |> last() |> map(fn: (r) => ({ r with _value: r._value - tables |> first() |> get(column: "_value") })) ) |> set(key: "_field", value: "volume") union(tables: [max, min, open, close, volume]) |> pivot(rowKey: ["_start"], columnKey: ["_field"], valueColumn: "_value") |> drop(columns: ["_stop", "_time"]) // 可选:移除不需要的字段
补充说明
- 如果窗口内只有一条Tick数据,计算出的
volume会是0,符合单条数据无新增成交量的逻辑。 - 若存在跨交易日的情况,需要额外过滤交易日边界,避免累计成交量重置导致的负数,可以通过添加时间范围过滤或结合交易日字段处理。
内容的提问来源于stack exchange,提问作者Heo
相关产品推荐
相关产品推荐

