InfluxDB使用full outer join报错:返回值被识别为函数
Flux全外连接(Full Outer Join)报错解决方法
我用以下Flux查询从InfluxDB获取汇总数据:
from(bucket: "external") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r["_measurement"] == "Elektro") |> filter(fn: (r) => r["_field"] == "Verbrauch") |> filter(fn: (r) => contains(value: r["_zaehler"], set: strings.split(v: "IT A1,IT A2,IT A3,Carrier A,IT B1,IT B2,IT B3,Carrier B", t: ","))) |> group(columns: ["_time"]) |> sum() |> map(fn: (r) => ({ r with _field: "Summe IT" })) |> group(columns: ["_time"])
得到的数据表包含_time和_value等字段(对应时间与IT总消耗值)。
尝试用全外连接将该表与自身按_time关联时,使用了如下查询:
join.full( left: ( from(bucket: "external") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r["_measurement"] == "Elektro") |> filter(fn: (r) => r["_field"] == "Verbrauch") |> filter(fn: (r) => contains(value: r["_zaehler"], set: strings.split(v: "IT A1,IT A2,IT A3,Carrier A,IT B1,IT B2,IT B3,Carrier B", t: ","))) |> group(columns: ["_time"]) |> sum() |> map(fn: (r) => ({ r with _field: "Summe IT" })) |> group(columns: ["_time"]) |> yield(name: "left") ), right: ( from(bucket: "external") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r["_measurement"] == "Elektro") |> filter(fn: (r) => r["_field"] == "Verbrauch") |> filter(fn: (r) => contains(value: r["_zaehler"], set: strings.split(v: "IT A1,IT A2,IT A3,Carrier A,IT B1,IT B2,IT B3,Carrier B", t: ","))) |> group(columns: ["_time"]) |> sum() |> map(fn: (r) => ({ r with _field: "Summe IT" })) |> group(columns: ["_time"]) |> yield(name: "right") ), on: (l, r) => (l._time == r._time), as: (l, r) => ({l}), )
执行后出现错误:
error @21:1-21:5: expected { A with full: ( as: (l: B, r: C) => {l: B}, left: stream[{D with _field: string}], on: (l: {E with _time: F}, r: {G with _time: H}) => bool, right: stream[{I with _field: string}], ) => stream[J], } (record) but found (<-tables: K, ?method: string, ?on: [string]) => stream[L] (function)
错误原因
- 误用join语法:Flux中没有
join.full()这种直接调用方式,全外连接需要通过join()函数并指定method: "full"参数来实现。 - 冗余的yield语句:
yield会生成额外的数据流输出,导致join无法获取到正确的单一输入流,必须移除。
修正后的查询
首先将重复的查询逻辑提取为变量,再使用正确的join语法:
// 定义通用的查询逻辑,方便后续替换_zaehler的筛选条件 baseQuery = () => from(bucket: "external") |> range(start: v.timeRangeStart, stop: v.timeRangeStop) |> filter(fn: (r) => r["_measurement"] == "Elektro") |> filter(fn: (r) => r["_field"] == "Verbrauch") |> filter(fn: (r) => contains(value: r["_zaehler"], set: strings.split(v: "IT A1,IT A2,IT A3,Carrier A,IT B1,IT B2,IT B3,Carrier B", t: ","))) |> group(columns: ["_time"]) |> sum() |> map(fn: (r) => ({ r with _field: "Summe IT" })) |> group(columns: ["_time"]) // 执行全外连接 join( method: "full", left: baseQuery(), right: baseQuery(), on: (l, r) => l._time == r._time, as: (l, r) => ({ l }) )
关键修改点
- 用
join(method: "full")替代join.full(),符合Flux的join函数规范。 - 移除left和right流中的
yield语句,确保输入为单一数据流。 - 将重复查询封装为函数
baseQuery,便于后续修改_zaehler的筛选条件,提升代码可维护性。
内容的提问来源于stack exchange,提问作者SirGamsay
相关产品推荐
相关产品推荐

