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

如何使用Flux实现时序表插值并关联两张数据源表?

问题描述

我有两张源表:
Table1(传感器数据):

time  | sensor
2     | 1.0
5     | 2.0
7     | 4.0

Table2(温度数据):

time  | temperature
1     | 20.0
10    | 30.0

时间精度为纳秒,需保留该精度。期望得到的结果表完全沿用Table1的时间序列,添加从Table2插值得到的温度值,最终结果如下:

Table Result
time  | sensor | temperature
2     | 1.0    | x 
5     | 2.0    | x
7     | 4.0    | x

我试过interpolate.linear(但它使用固定间隔),也试过fill函数(但它不支持插值),请问如何通过插值关联两张表?

解决方案

可以通过Flux的join函数结合自定义线性插值逻辑实现,以下分两种场景给出实现代码:

通用场景(Table2有多个时间点)

适用于Table2包含任意数量时间点的情况,自动匹配每个Table1时间点前后的温度记录进行插值:

// 读取并整理Table1数据
table1 = from(bucket: "your-bucket")
  |> range(start: -inf)
  |> filter(fn: (r) => r._measurement == "Table1" and r._field == "sensor")
  |> rename(columns: {_value: "sensor"})
  |> keep(columns: ["_time", "sensor"])

// 读取Table2数据并按时间排序
table2 = from(bucket: "your-bucket")
  |> range(start: -inf)
  |> filter(fn: (r) => r._measurement == "Table2" and r._field == "temperature")
  |> rename(columns: {_value: "temperature"})
  |> keep(columns: ["_time", "temperature"])
  |> sort(columns: ["_time"])

// 自定义线性插值逻辑
linearInterpolate = (tables=<-) =>
  tables
    // 关联Table1与Table2的所有记录
    |> join(
      tables: {t1: table1, t2: table2},
      on: [],
      fn: (t1, t2) => ({t1 with t2_temp: t2.temperature, t2_time: t2._time})
    )
    // 筛选出Table1时间点之前的温度记录
    |> filter(fn: (r) => r._time >= r.t2_time)
    // 按Table1的时间和传感器值分组
    |> group(columns: ["_time", "sensor"])
    |> sort(columns: ["t2_time"])
    // 全局窗口,便于获取前后时间点
    |> window(every: inf)
    // 提取每个时间点的前后温度记录
    |> reduce(
      fn: (r, accumulator) => ({
        prev_time: r.t2_time,
        prev_temp: r.t2_temp,
        next_time: if r.t2_time > accumulator.prev_time then r.t2_time else accumulator.next_time,
        next_temp: if r.t2_time > accumulator.prev_time then r.t2_temp else accumulator.next_temp,
        sensor: r.sensor,
        _time: r._time
      }),
      identity: {prev_time: 0, prev_temp: 0.0, next_time: 0, next_temp: 0.0, sensor: 0.0, _time: 0}
    )
    // 计算线性插值温度
    |> map(
      fn: (r) => ({
        _time: r._time,
        sensor: r.sensor,
        temperature: r.prev_temp + (r.next_temp - r.prev_temp) * float(v: r._time - r.prev_time) / float(v: r.next_time - r.prev_time)
      })
    )
    |> keep(columns: ["_time", "sensor", "temperature"])

// 执行插值计算
linearInterpolate()

简化场景(Table2仅含固定插值点)

如果Table2只有示例中的两个固定时间点,可以直接提取这两个点计算插值,效率更高:

// 读取并整理Table1数据
table1 = from(bucket: "your-bucket")
  |> range(start: -inf)
  |> filter(fn: (r) => r._measurement == "Table1" and r._field == "sensor")
  |> rename(columns: {_value: "sensor"})
  |> keep(columns: ["_time", "sensor"])

// 获取Table2的两个固定时间点数据
temp_points = from(bucket: "your-bucket")
  |> range(start: -inf)
  |> filter(fn: (r) => r._measurement == "Table2" and r._field == "temperature")
  |> rename(columns: {_value: "temperature"})
  |> keep(columns: ["_time", "temperature"])
  |> sort(columns: ["_time"])
  |> findRecord(fn: (key) => true, idx: 0)
  |> (fn: (first) => 
    from(bucket: "your-bucket")
      |> range(start: -inf)
      |> filter(fn: (r) => r._measurement == "Table2" and r._field == "temperature")
      |> rename(columns: {_value: "temperature"})
      |> keep(columns: ["_time", "temperature"])
      |> sort(columns: ["_time"])
      |> findRecord(fn: (key) => true, idx: 1)
      |> (fn: (second) => ({first: first, second: second}))
  )

// 直接计算每个Table1时间点的插值温度
table1
  |> map(
    fn: (r) => ({
      _time: r._time,
      sensor: r.sensor,
      temperature: temp_points.first.temperature + (temp_points.second.temperature - temp_points.first.temperature) * float(v: r._time - temp_points.first._time) / float(v: temp_points.second._time - temp_points.first._time)
    })
  )

说明

两种方案均保留纳秒级时间精度,通过线性插值公式计算对应温度值,完全贴合需求中沿用Table1时间序列的要求。

内容的提问来源于stack exchange,提问作者Pablo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 22:25:18