InfluxDB按条件提取周期数据的查询实现求助
实现InfluxDB基于条件的周期数据查询
针对你需要提取wp1或wp2完整周期数据的需求,我会分**InfluxDB 2.x(Flux查询语言)和InfluxDB 1.x(InfluxQL)**两种场景给出解决方案,先明确周期的判定规则:
- 周期起始:字段值首次从≤0.5变为>0.5的行
- 周期中:字段值持续>0.5的所有行
- 周期结束:字段值首次从>0.5变为≤0.5的行
一、InfluxDB 2.x 使用Flux查询
Flux的状态函数非常适合处理这种依赖前后行状态变化的场景,以下是具体实现:
1. 提取单个字段(如wp1)的完整周期数据
from(bucket: "你的桶名") |> range(start: -inf, stop: +inf) |> filter(fn: (r) => r._measurement == "Power") |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value") |> sort(columns: ["_time"]) |> statefulMap(fn: (r, state) => { // 初始化状态:记录前一行的wp1值和是否处于周期中 prev_wp1 = if exists state.prev_wp1 then state.prev_wp1 else -1.0 in_cycle = if exists state.in_cycle then state.in_cycle else false // 标记当前行的类型 is_start = not in_cycle and r.wp1 > 0.5 is_end = in_cycle and r.wp1 <= 0.5 is_mid = in_cycle and r.wp1 > 0.5 // 更新状态 new_state = {prev_wp1: r.wp1, in_cycle: (in_cycle or is_start) and not is_end} // 返回带标记的行 return ({r with is_keep: is_start or is_mid or is_end}, new_state) }) |> filter(fn: (r) => r.is_keep) |> keep(columns: ["_time", "id", "level", "wp1", "wp2", "wp3"])
如果只需要每个周期的起始行、最后一个周期中行、结束行(和你给出的示例一致),可以在上述查询后追加:
|> group(columns: ["id"]) |> window(every: inf) |> map(fn: (r) => { start_rows = filter(fn: (x) => x.is_start, arr: r._values) mid_rows = filter(fn: (x) => x.is_mid, arr: r._values) end_rows = filter(fn: (x) => x.is_end, arr: r._values) result = start_rows if length(arr: mid_rows) > 0 then result = result ++ [last(arr: mid_rows)] else result result = result ++ end_rows return result }) |> flatten()
2. 同时提取wp1和wp2的周期数据
只需修改状态跟踪逻辑,同时监控两个字段的状态变化:
from(bucket: "你的桶名") |> range(start: -inf, stop: +inf) |> filter(fn: (r) => r._measurement == "Power") |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value") |> sort(columns: ["_time"]) |> statefulMap(fn: (r, state) => { // 分别跟踪wp1和wp2的状态 prev_wp1 = if exists state.prev_wp1 then state.prev_wp1 else -1.0 prev_wp2 = if exists state.prev_wp2 then state.prev_wp2 else -1.0 in_cycle_wp1 = if exists state.in_cycle_wp1 then state.in_cycle_wp1 else false in_cycle_wp2 = if exists state.in_cycle_wp2 then state.in_cycle_wp2 else false // 标记wp1的行类型 is_start_wp1 = not in_cycle_wp1 and r.wp1 > 0.5 is_end_wp1 = in_cycle_wp1 and r.wp1 <= 0.5 is_mid_wp1 = in_cycle_wp1 and r.wp1 > 0.5 // 标记wp2的行类型 is_start_wp2 = not in_cycle_wp2 and r.wp2 > 0.5 is_end_wp2 = in_cycle_wp2 and r.wp2 <= 0.5 is_mid_wp2 = in_cycle_wp2 and r.wp2 > 0.5 // 更新状态 new_state = { prev_wp1: r.wp1, prev_wp2: r.wp2, in_cycle_wp1: (in_cycle_wp1 or is_start_wp1) and not is_end_wp1, in_cycle_wp2: (in_cycle_wp2 or is_start_wp2) and not is_end_wp2 } // 只要属于任意一个字段的周期,就保留该行 keep_row = is_start_wp1 or is_mid_wp1 or is_end_wp1 or is_start_wp2 or is_mid_wp2 or is_end_wp2 return ({r with keep_row: keep_row}, new_state) }) |> filter(fn: (r) => r.keep_row) |> keep(columns: ["_time", "id", "level", "wp1", "wp2", "wp3"])
二、InfluxDB 1.x 使用InfluxQL查询
InfluxQL没有原生的状态函数,需要借助LAG()窗口函数来比较前后行的值:
1. 提取单个字段(如wp1)的完整周期数据
SELECT time, id, level, wp1, wp2, wp3 FROM ( SELECT time, id, level, wp1, wp2, wp3, LAG(wp1) OVER (ORDER BY time) AS prev_wp1 FROM Power ) t WHERE -- 周期起始:当前>0.5,前一行≤0.5(或无前行) (wp1 > 0.5 AND (prev_wp1 IS NULL OR prev_wp1 <= 0.5)) -- 周期中:当前和前一行都>0.5 OR (wp1 > 0.5 AND prev_wp1 > 0.5) -- 周期结束:当前≤0.5,前一行>0.5 OR (wp1 <= 0.5 AND prev_wp1 > 0.5) ORDER BY time;
如果需要每个周期的起始行、最后一个周期中行、结束行,可以用以下嵌套查询:
WITH cycle_boundaries AS ( SELECT CASE WHEN wp1 > 0.5 AND (prev_wp1 IS NULL OR prev_wp1 <= 0.5) THEN time END AS start_time, CASE WHEN wp1 <= 0.5 AND prev_wp1 > 0.5 THEN time END AS end_time FROM ( SELECT time, wp1, LAG(wp1) OVER (ORDER BY time) AS prev_wp1 FROM Power ) t ), paired_cycles AS ( SELECT start_time, (SELECT MIN(end_time) FROM cycle_boundaries WHERE end_time > start_time) AS end_time FROM cycle_boundaries WHERE start_time IS NOT NULL ) SELECT * FROM Power WHERE time IN (SELECT start_time FROM paired_cycles) OR time IN (SELECT end_time FROM paired_cycles) OR time IN ( SELECT MAX(time) FROM Power WHERE wp1 > 0.5 GROUP BY ( SELECT start_time FROM paired_cycles WHERE start_time <= time AND end_time >= time ) ) ORDER BY time;
2. 同时提取wp1和wp2的周期数据
用UNION ALL合并两个字段的查询结果:
-- wp1的周期数据 SELECT time, id, level, wp1, wp2, wp3, 'wp1' AS cycle_source FROM ( SELECT time, id, level, wp1, wp2, wp3, LAG(wp1) OVER (ORDER BY time) AS prev_wp1 FROM Power ) t WHERE (wp1 > 0.5 AND (prev_wp1 IS NULL OR prev_wp1 <= 0.5)) OR (wp1 > 0.5 AND prev_wp1 > 0.5) OR (wp1 <= 0.5 AND prev_wp1 > 0.5) UNION ALL -- wp2的周期数据 SELECT time, id, level, wp1, wp2, wp3, 'wp2' AS cycle_source FROM ( SELECT time, id, level, wp1, wp2, wp3, LAG(wp2) OVER (ORDER BY time) AS prev_wp2 FROM Power ) t WHERE (wp2 > 0.5 AND (prev_wp2 IS NULL OR prev_wp2 <= 0.5)) OR (wp2 > 0.5 AND prev_wp2 > 0.5) OR (wp2 <= 0.5 AND prev_wp2 > 0.5) ORDER BY time;
内容的提问来源于stack exchange,提问作者Md Enayat
相关产品推荐
相关产品推荐

