You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 14:57:33