如何在Go中解析Influx查询结果并转换为结构体数组?
如何将InfluxDB查询结果按时间分组并映射为结构体数组
结构体定义
type CandlesticksHistoricalData struct { OpenTime time.Time Open decimal.Decimal High decimal.Decimal Low decimal.Decimal // 可选:若需收盘价,可添加Close字段 // Close decimal.Decimal } var result []CandlesticksHistoricalData
Flux查询语句
from(bucket: "DataManager") |> range(start: 1638316800, stop:1638316980) |> filter(fn: (r) => r._measurement == "Candlesticks") |> filter(fn: (r) => r.exchange == "Deribit") |> filter(fn: (r) => r.symbol == "BTCUSD") |> filter(fn: (r) => r.tf == "1Min") |> filter(fn: (r) => r._field == "o" or r._field == "h" or r._field == "l" or r._field == "c")
当前问题
查询返回的结果按字段分组(先返回所有时间的c字段,再返回h字段,以此类推),无法直接映射到结构体。需要将同时间的o、h、l、c字段值聚合到同一个结构体中,并按时间顺序整理结果。
当前代码输出日志:
Time: 2021-12-01 00:00:00 +0000 UTC, Value : 56922.5, Field c, TableChanged: true Time: 2021-12-01 00:01:00 +0000 UTC, Value : 56802, Field c, TableChanged: false Time: 2021-12-01 00:02:00 +0000 UTC, Value : 56730.5, Field c, TableChanged: false Time: 2021-12-01 00:00:00 +0000 UTC, Value : 57073, Field h, TableChanged: false Time: 2021-12-01 00:01:00 +0000 UTC, Value : 56946.5, Field h, TableChanged: false Time: 2021-12-01 00:02:00 +0000 UTC, Value : 56822, Field h, TableChanged: false Time: 2021-12-01 00:00:00 +0000 UTC, Value : 56897, Field l, TableChanged: false Time: 2021-12-01 00:01:00 +0000 UTC, Value : 56800, Field l, TableChanged: false Time: 2021-12-01 00:02:00 +0000 UTC, Value : 56730.5, Field l, TableChanged: false Time: 2021-12-01 00:00:00 +0000 UTC, Value : 56995.5, Field o, TableChanged: false Time: 2021-12-01 00:01:00 +0000 UTC, Value : 56907.5, Field o, TableChanged: false Time: 2021-12-01 00:02:00 +0000 UTC, Value : 56810.5, Field o, TableChanged: false
期望按时间分组输出,同时间的所有字段连续出现,最终映射为结构体数组。
解决方案
方案一:客户端侧分组处理
通过临时Map按时间分组,收集同时间的所有字段值,最后整理为结构体数组并排序:
result, err := dm.fluxConn.QueryApi.Query(context.Background(), fluxQuery) if err != nil { return nil, err } defer result.Close() // 临时存储按时间分组的K线数据 tempMap := make(map[time.Time]*CandlesticksHistoricalData) for result.Next() { record := result.Record() t := record.Time() v := record.Value() f := record.Field() // 将查询值转换为decimal.Decimal val, ok := v.(float64) if !ok { // 可根据需求处理类型转换错误,比如跳过或报错 continue } decVal := decimal.NewFromFloat(val) // 初始化或获取当前时间对应的结构体 cs, exists := tempMap[t] if !exists { cs = &CandlesticksHistoricalData{ OpenTime: t, } tempMap[t] = cs } // 根据字段名赋值 switch f { case "o": cs.Open = decVal case "h": cs.High = decVal case "l": cs.Low = decVal case "c": // 若结构体定义了Close字段,取消注释以下代码 // cs.Close = decVal } } // 检查查询过程中是否出现错误 if err := result.Err(); err != nil { return nil, err } // 将Map中的数据转换为数组并按时间排序 var resultSlice []CandlesticksHistoricalData for _, cs := range tempMap { resultSlice = append(resultSlice, *cs) } sort.Slice(resultSlice, func(i, j int) bool { return resultSlice[i].OpenTime.Before(resultSlice[j].OpenTime) }) return resultSlice, nil
方案二:修改Flux查询,返回宽表数据
使用Flux的pivot函数将字段转为列,让每条记录直接包含完整的K线数据,避免客户端分组处理:
修改后的Flux查询
from(bucket: "DataManager") |> range(start: 1638316800, stop:1638316980) |> filter(fn: (r) => r._measurement == "Candlesticks") |> filter(fn: (r) => r.exchange == "Deribit") |> filter(fn: (r) => r.symbol == "BTCUSD") |> filter(fn: (r) => r.tf == "1Min") |> filter(fn: (r) => r._field == "o" or r._field == "h" or r._field == "l" or r._field == "c") |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
对应客户端代码
result, err := dm.fluxConn.QueryApi.Query(context.Background(), fluxQuery) if err != nil { return nil, err } defer result.Close() var resultSlice []CandlesticksHistoricalData for result.Next() { record := result.Record() t := record.Time() // 获取各字段值并转换类型 oVal, oOk := record.ValueByKey("o").(float64) hVal, hOk := record.ValueByKey("h").(float64) lVal, lOk := record.ValueByKey("l").(float64) // cVal, cOk := record.ValueByKey("c").(float64) // 检查字段是否存在且类型正确 if !oOk || !hOk || !lOk { continue } // 构建结构体并添加到数组 cs := CandlesticksHistoricalData{ OpenTime: t, Open: decimal.NewFromFloat(oVal), High: decimal.NewFromFloat(hVal), Low: decimal.NewFromFloat(lVal), // Close: decimal.NewFromFloat(cVal), // 若需收盘价 } resultSlice = append(resultSlice, cs) } if err := result.Err(); err != nil { return nil, err } // 可选:若查询结果未按时间排序,执行排序 sort.Slice(resultSlice, func(i, j int) bool { return resultSlice[i].OpenTime.Before(resultSlice[j].OpenTime) }) return resultSlice, nil
两种方案中,方案二更高效,因为将分组逻辑交给InfluxDB处理,减少客户端的代码复杂度和数据处理量。
内容的提问来源于stack exchange,提问作者Lev Siemens
相关产品推荐
相关产品推荐

