InfluxDB v3多表Join时单表无结果如何保留有效数据?
问题描述
我用InfluxDB查询监控指标,写了三个查询分别提取memory、disk、cpu的数据,想用join()函数把三张表合并成一张。但如果其中某张表(比如这次的cpu表)没有查询结果,合并后的整张表也不会返回任何结果。我希望就算部分表无数据,也能返回其他有数据的表的内容。以下是我的Flux查询代码:
//Define your queries for different measurements memory=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "mem") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) disk=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "disk") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) cpu=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "cpu") |> filter(fn: (r) => r["_field"] == "usage_idle") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value": float(v: r["_value"])})) |> map(fn: (r) => ({ r with _value: 100.0 - r._value })) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) result=join(tables:{memory:memory, disk:disk}, on:["_time"]) join(tables: {result: result,cpu: cpu}, on:["_time"]) |>yield()
解决方案
Flux的join()默认是内连接,只有当所有参与连接的表在on指定的字段(这里是_time)上都有匹配值时,才会保留该行数据。如果某张表为空,内连接的结果自然为空。要实现“部分表无数据仍返回其他有数据内容”的需求,可采用以下两种方案:
方案1:使用外连接(Outer Join)
通过join()的method参数指定为"outer",保留所有表中存在的_time行,不存在的字段值填充为null。
修改后的完整查询代码:
// 定义各指标查询 memory=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "mem") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) disk=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "disk") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) cpu=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "cpu") |> filter(fn: (r) => r["_field"] == "usage_idle") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value": float(v: r["_value"])})) |> map(fn: (r) => ({ r with _value: 100.0 - r._value })) |> map(fn:(r) => ({r with "_value" : string(v: r["_value"])})) |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value") |> group() |> sort(columns: ["_time"], desc: false) // 第一步外连接memory和disk result = join( tables: {memory: memory, disk: disk}, on: ["_time"], method: "outer" ) // 第二步将结果与cpu做外连接 join( tables: {result: result, cpu: cpu}, on: ["_time"], method: "outer" ) |> yield()
方案2:用union合并后重新整理结构
如果外连接的结果不符合预期,也可以先将三个表用union合并,再通过pivot重新按_time整理列。这种方式会自动保留所有存在的_time行,缺失的字段值为null。
修改后的代码:
// 定义各指标查询,去掉pivot,重命名字段便于后续整理 memory=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "mem") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({ _time: r._time, host: r.host, mem_used: string(v: r._value) })) |> group() |> sort(columns: ["_time"], desc: false) disk=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "disk") |> filter(fn: (r) => r["_field"] == "used") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({ _time: r._time, host: r.host, disk_used: string(v: r._value) })) |> group() |> sort(columns: ["_time"], desc: false) cpu=from(bucket: "telegraf") |> range(start : -2d) |> filter(fn: (r) => r["_measurement"] == "cpu") |> filter(fn: (r) => r["_field"] == "usage_idle") |> filter(fn: (r) => r["host"] == "host") |> limit(n:20) |> map(fn:(r) => ({r with "_value": float(v: r["_value"])})) |> map(fn: (r) => ({ r with _value: 100.0 - r._value })) |> map(fn:(r) => ({ _time: r._time, host: r.host, cpu_used: string(v: r._value) })) |> group() |> sort(columns: ["_time"], desc: false) // 合并三个表 union(tables: [memory, disk, cpu]) // 按时间和主机分组,确保同时间同主机的指标合并 |> group(columns: ["_time", "host"]) // 转置为宽表,缺失字段显示null |> pivot(rowKey: ["_time", "host"], columnKey: ["_field"], valueColumn: "_value") |> yield()
关键说明
- 外连接(
method: "outer")适合需要严格按时间对齐的场景,会保留所有参与连接表中的时间记录,缺失字段填充null。 union+pivot的方式更灵活,无需提前做pivot处理,合并后直接整理结构,同样能保留所有存在的时间点数据。
内容的提问来源于stack exchange,提问作者khaled
相关产品推荐
相关产品推荐

