如何使用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
相关产品推荐
相关产品推荐

