Skip to content

Flux 查询语法

Flux 是一种面向数据流的查询和处理语言,主要用于 InfluxDB 2.x。它不像 SQL 那样围绕一张二维表编写 SELECT,而是把数据作为“表流”依次传给过滤、分组、变换和聚合函数。

InfluxDB 3 不支持 Flux。新建 InfluxDB 3 系统应使用 SQL 或 InfluxQL;本篇适用于维护 InfluxDB 2.x、已有 Dashboard 和 Task,或阅读旧项目中的 Flux 查询。


最小查询

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")

执行过程:

text
读取 iot Bucket
  -> 保留最近 1 小时
  -> 保留 sensor_data Measurement
  -> 保留 temperature Field
  -> 返回符合条件的表流

|> 称为 Pipe-forward Operator。它把左侧结果作为第一个参数传给右侧函数:

flux
data |> mean()

可以理解成:

text
mean(tables: data)

Flux 数据结构

Line Protocol 中的一条 Point:

text
sensor_data,site=hefei,device_id=sensor01 temperature=26.3,humidity=61.5 1787212800000000000

使用 from() 查询后,Field 默认采用“窄表”形式,每个 Field Value 是独立记录:

_time_measurementsitedevice_id_field_value
08:00sensor_datahefeisensor01temperature26.3
08:00sensor_datahefeisensor01humidity61.5

核心列:

含义
_timePoint 的 Timestamp
_measurementMeasurement 名称
_fieldField Key,例如 temperature
_value对应 Field Value,例如 26.3
_start当前查询时间范围起点
_stop当前查询时间范围终点
其他列Tag,例如 sitedevice_id

这与 InfluxDB 3 SQL 直接看到 temperaturehumidity 两列不同。Flux 中经常先按 _field 过滤;需要把多个 Field 变回普通列时,使用 pivot()


变量与基本类型

变量

flux
bucketName = "iot"
queryStart = -24h
targetSite = "hefei"

from(bucket: bucketName)
    |> range(start: queryStart)
    |> filter(fn: (r) => r.site == targetSite)

Flux 变量不可重新赋值。变量命名区分大小写。

常用字面量

flux
stringValue = "hefei"
integerValue = 10
floatValue = 26.3
booleanValue = true
durationValue = 5m
timeValue = 2026-08-20T08:00:00Z
arrayValue = ["hefei", "shanghai"]

常见 Duration 单位包括 nsusmssmhdw

flux
range(start: -15m)
aggregateWindow(every: 1h, fn: mean)

Record

Flux 中一行数据称为 Record,通常用 r 表示:

flux
filter(fn: (r) => r.site == "hefei")

列名包含特殊字符时使用方括号:

flux
filter(fn: (r) => r["device-id"] == "sensor01")

from 与 range

from() 从 Bucket 读取数据:

flux
from(bucket: "iot")

从 InfluxDB 读取数据后必须通过 range() 指定时间范围:

flux
// 最近 1 小时
from(bucket: "iot")
    |> range(start: -1h)
flux
// 固定时间范围,左闭右开
from(bucket: "iot")
    |> range(
        start: 2026-08-20T08:00:00Z,
        stop: 2026-08-20T09:00:00Z,
    )
flux
// 使用 Dashboard 或 API 提供的时间范围
from(bucket: "iot")
    |> range(start: v.timeRangeStart, stop: v.timeRangeStop)

range() 不只是普通过滤,它还设置结果的 _start_stop,这些列通常属于 Group Key。


filter 条件过滤

Measurement、Field 与 Tag

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> filter(fn: (r) => r.site == "hefei")
    |> filter(fn: (r) => r.device_id == "sensor01")

也可以合并为一个条件:

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) =>
        r._measurement == "sensor_data" and
        r._field == "temperature" and
        r.site == "hefei" and
        r.device_id == "sensor01"
    )

多值条件

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) =>
        r._field == "temperature" or r._field == "humidity"
    )

使用 contains() 判断数组是否包含某个值:

flux
targetDevices = ["sensor01", "sensor02"]

from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => contains(value: r.device_id, set: targetDevices))

contains() 对大数组可能效率较低。固定的少量值可直接使用 or;大规模动态集合应重新评估 Schema 或查询方案。

正则表达式

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => r.device_id =~ /^sensor0[1-3]$/)

=~ 表示匹配,!~ 表示不匹配。能用精确比较时优先使用 ==,正则通常需要更多计算。

数值条件

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._field == "temperature")
    |> filter(fn: (r) => r._value >= 30.0)

先限定 _field 再比较 _value,可避免不同 Field 类型混在一起造成类型错误或错误语义。


keep、drop 与 rename

只保留需要的列:

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> keep(columns: ["_time", "_field", "_value", "site", "device_id"])

删除指定列:

flux
data
    |> drop(columns: ["_start", "_stop"])

重命名列:

flux
data
    |> rename(columns: {_value: "temperature"})

keep()drop() 若移除 Group Key 中的列,会改变输出表的分组结构。


sort、limit 与 tail

按时间升序:

flux
data
    |> sort(columns: ["_time"])

按时间降序并取前 10 行:

flux
data
    |> sort(columns: ["_time"], desc: true)
    |> limit(n: 10)

取每张输入表的最后 10 行:

flux
data
    |> tail(n: 10)

这些函数通常对表流中的每张 Table 分别执行,不一定是对整个 Bucket 的全局结果执行。是否需要先 group() 合并表,取决于希望得到“每台设备 10 条”还是“所有设备一共 10 条”。


聚合函数

常用聚合和 Selector:

函数作用
count()统计记录数
mean()平均值
sum()求和
min()max()最小值、最大值
first()last()最早值、最新值
median()中位数
quantile()分位数
stddev()标准差

查询每组温度的平均值:

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> mean()

mean() 对每张输入 Table 分别计算。默认 Group Key 通常包含 Measurement、Field 和 Tag,因此结果可能是每个设备各一行,而不是所有设备合成一行。


group 与 Group Key

Flux 查询返回的是一组 Table。每张 Table 都有 Group Key,Group Key 中各列在该 Table 内具有相同值。

例如默认可能按以下列分表:

text
[_start, _stop, _field, _measurement, site, device_id]

按站点重新分组:

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> group(columns: ["site"])
    |> mean()

此时得到每个站点的平均温度,而不是每台设备的平均温度。

把所有输入合成一个逻辑分组:

flux
data
    |> group(columns: [])
    |> mean()

取消分组会混合不同设备甚至不同 Field,使用前必须确认业务语义。Flux 中许多“结果为什么有很多张表”的问题,都与 Group Key 有关。


aggregateWindow 时间窗口

按 5 分钟计算平均温度:

flux
from(bucket: "iot")
    |> range(start: -24h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> aggregateWindow(every: 5m, fn: mean, createEmpty: false)

关键参数:

参数含义
every窗口间隔
fn每个窗口执行的函数
createEmpty是否创建没有数据的空窗口
offset调整窗口边界
timeSrc从哪一列获取输出时间,默认常用 _stop
timeDst输出时间列,通常为 _time

按 1 小时取最大值:

flux
data
    |> aggregateWindow(every: 1h, fn: max)

aggregateWindow() 通常比手工组合 window()mean()duplicate() 更适合 Dashboard 降采样。

Dashboard 可以使用自动窗口变量:

flux
from(bucket: "iot")
    |> range(start: v.timeRangeStart, stop: v.timeRangeStop)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> aggregateWindow(every: v.windowPeriod, fn: mean)

pivot 将 Field 转成列

查询温度和湿度时,默认结果使用 _field_value 两列:

flux
data = from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) =>
        r._field == "temperature" or r._field == "humidity"
    )

使用 pivot() 转为宽表:

flux
data
    |> pivot(
        rowKey: ["_time"],
        columnKey: ["_field"],
        valueColumn: "_value",
    )
    |> keep(columns: ["_time", "site", "device_id", "temperature", "humidity"])

结果:

_timesitedevice_idtemperaturehumidity
08:00hefeisensor0126.361.5

当后续计算需要同时访问多个 Field 时,通常先 pivot()

flux
data
    |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value")
    |> filter(fn: (r) => r.temperature >= 30.0 and r.humidity >= 60.0)

pivot() 可能显著增加内存使用。应先通过 range()、Measurement、Field 和 Tag 缩小数据,再 Pivot。


map 计算新列

将摄氏温度转换为华氏温度:

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> map(fn: (r) => ({
        r with
        _value: r._value * 9.0 / 5.0 + 32.0,
        _field: "temperature_f",
    }))

{r with ...} 表示保留原 Record 的其他列,只覆盖或增加指定列。若直接返回 {_time: r._time, _value: ...},未显式保留的 Tag 和其他列会丢失。

增加告警级别:

flux
data
    |> map(fn: (r) => ({
        r with
        level: if r._value >= 35.0 then "critical"
            else if r._value >= 30.0 then "warning"
            else "normal",
    }))

Flux 条件表达式必须包含 else,并产生一个值。


类型转换

常用转换函数:

flux
data |> toFloat()
data |> toInt()
data |> toString()
data |> toBool()

map() 中转换某列:

flux
data
    |> map(fn: (r) => ({r with retry_count: int(v: r.retry_count)}))

转换不能修复已经写入的 Field 类型冲突。更合理的做法是在采集和写入阶段保持类型稳定。


difference、derivative 与 increase

相邻值差

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "power_meter")
    |> filter(fn: (r) => r._field == "total_kwh")
    |> difference(nonNegative: true)

difference() 计算相邻记录的值差。nonNegative: true 可把负差视为异常重置场景,但是否符合计数器语义需要结合数据判断。

单位时间变化率

flux
data
    |> derivative(unit: 1m, nonNegative: true)

derivative() 根据值差和时间差计算变化率。

计数器增长量

flux
data
    |> increase()

increase() 常用于可能重置的单调计数器。温度等 Gauge 不适合使用计数器增长函数。


fill 与缺失值

使用固定值替换 null

flux
data
    |> fill(column: "_value", value: 0.0)

使用上一条非空值:

flux
data
    |> fill(column: "_value", usePrevious: true)

若希望为空时间段创建窗口,需要聚合时启用空窗口:

flux
data
    |> aggregateWindow(every: 5m, fn: mean, createEmpty: true)
    |> fill(column: "_value", usePrevious: true)

填充值不是真实采样。设备离线时填 0、沿用上一值或保留空值代表不同业务含义,不能只为了让曲线连续而随意选择。


union 与 join

union 合并相同结构的数据流

flux
current = from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")

history = from(bucket: "iot_archive")
    |> range(start: -24h, stop: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")

union(tables: [history, current])
    |> sort(columns: ["_time"])

union() 合并表流,但不保证结果顺序,因此需要时应再 sort()

join 关联两个数据流

将温度与功率按设备和时间关联:

flux
temperature = from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> aggregateWindow(every: 5m, fn: mean)
    |> rename(columns: {_value: "temperature"})

power = from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "power_data")
    |> filter(fn: (r) => r._field == "watts")
    |> aggregateWindow(every: 5m, fn: mean)
    |> rename(columns: {_value: "watts"})

join(
    tables: {temperature: temperature, power: power},
    on: ["_time", "device_id"],
)
    |> keep(columns: ["_time", "device_id", "temperature", "watts"])

两侧原始 Timestamp 经常无法完全相等,因此示例先聚合到相同的 5 分钟窗口。JOIN 前应分别限制时间范围、过滤 Field,并确保关联键符合预期,否则会产生大量中间数据。


exists 与空列判断

判断 Record 是否包含某列且值非空:

flux
data
    |> filter(fn: (r) => exists r.temperature)

常与 pivot() 一起使用:

flux
data
    |> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value")
    |> filter(fn: (r) => exists r.temperature and exists r.humidity)

yield 与多个结果

yield() 为输出结果命名:

flux
data
    |> mean()
    |> yield(name: "mean_temperature")

一个脚本可以输出多个结果:

flux
data = from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")

data |> mean() |> yield(name: "mean")
data |> max() |> yield(name: "max")

没有显式 yield() 时,脚本中未被继续消费的表流通常会隐式输出。Task 中使用 to() 写回数据时,一般不需要为同一分支再 yield()


package、import 与自定义函数

导入标准库包:

flux
import "math"

data
    |> map(fn: (r) => ({r with _value: math.round(x: r._value)}))

查看 Schema:

flux
import "influxdata/influxdb/schema"

schema.measurements(bucket: "iot")
flux
import "influxdata/influxdb/schema"

schema.measurementFieldKeys(
    bucket: "iot",
    measurement: "sensor_data",
)

定义可复用函数:

flux
temperatureData = (bucket, start, site) =>
    from(bucket: bucket)
        |> range(start: start)
        |> filter(fn: (r) => r._measurement == "sensor_data")
        |> filter(fn: (r) => r._field == "temperature")
        |> filter(fn: (r) => r.site == site)

temperatureData(bucket: "iot", start: -1h, site: "hefei")
    |> aggregateWindow(every: 5m, fn: mean)

to 写回 InfluxDB

将处理结果写入另一个 Bucket:

flux
from(bucket: "iot_raw")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> aggregateWindow(every: 5m, fn: mean)
    |> to(bucket: "iot_downsampled", org: "example-org")

to() 写入需要结果保留 _time_measurement_field_value,其他非系统列通常作为 Tag。前面的 keep()drop()map()pivot() 若破坏这些列,需要在写回前重新整理 Schema。

为了避免任务重复处理自己写出的数据,原始数据与降采样数据通常放在不同 Bucket,或使用不同 Measurement 并严格过滤输入。


Task 定时处理

Task 是定时执行的 Flux 脚本,常用于降采样、异常检测和数据转存。

每小时运行一次,把最近一小时数据按 5 分钟降采样:

flux
option task = {
    name: "sensor_downsample_5m",
    every: 1h,
    offset: 5m,
}

from(bucket: "iot_raw")
    |> range(start: -task.every)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) =>
        r._field == "temperature" or r._field == "humidity"
    )
    |> aggregateWindow(every: 5m, fn: mean)
    |> to(bucket: "iot_downsampled", org: "example-org")

Task 配置:

属性作用
nameTask 名称
every固定间隔执行
cron使用 Cron 表达式执行,与 every 二选一
offset延迟执行但保持原计划窗口,用于等待迟到数据

offset: 5m 表示整点任务延迟 5 分钟执行,使迟到数据有时间到达;它不是把查询窗口整体向后移动。

创建 Task:

bash
influx task create \
  --org example-org \
  --file sensor-downsample.flux

Task 需要考虑失败重跑、重复输出、迟到数据和窗口边界。仅使用 range(start: -task.every) 的任务若手动补跑,可能与原定窗口不同;重要场景应结合 option task 提供的运行时间边界设计。


SQL、InfluxQL 与 Flux 对照

最近一小时温度

Flux:

flux
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> filter(fn: (r) => r._field == "temperature")
    |> filter(fn: (r) => r.site == "hefei")

InfluxQL:

sql
SELECT "temperature"
FROM "sensor_data"
WHERE time >= now() - 1h
  AND "site" = 'hefei';

InfluxDB 3 SQL:

sql
SELECT time, site, device_id, temperature
FROM sensor_data
WHERE time >= now() - INTERVAL '1 hour'
  AND site = 'hefei';

按 5 分钟平均

语言时间窗口
FluxaggregateWindow(every: 5m, fn: mean)
InfluxQLGROUP BY time(5m) + MEAN()
InfluxDB 3 SQLdate_bin(INTERVAL '5 minutes', time) + avg()

常见错误

忘记 range

flux
// 错误:from 后没有时间范围
from(bucket: "iot")
    |> filter(fn: (r) => r._measurement == "sensor_data")

读取 InfluxDB 时应紧跟 range(),并限制到业务真正需要的范围。

忘记过滤 Field

flux
// temperature、humidity 等不同 Field 都会进入 mean
from(bucket: "iot")
    |> range(start: -1h)
    |> filter(fn: (r) => r._measurement == "sensor_data")
    |> mean()

若希望计算温度平均值,必须增加 _field == "temperature"

不理解默认分组

mean()last() 等函数对每张输入 Table 执行。结果比预期多或少时,检查 Group Key,并在聚合前显式 group(columns: [...])

过早 pivot

先 Pivot 全部历史数据再过滤,会制造庞大的中间结果。推荐顺序:

text
from -> range -> filter Measurement/Field/Tag -> aggregate -> pivot -> map

具体顺序仍取决于计算是否要求原始多字段对齐。

把 _value 当成固定类型

不同 _field_value 可能是 Float、Integer、String 或 Boolean。执行数值比较前先过滤 Field,避免把不同类型放入同一计算。

在 InfluxDB 3 中执行 Flux

InfluxDB 3 没有 Flux 查询引擎。迁移时不仅要替换 from(),还要重新处理窄表/宽表、Group Key、窗口边界、Task 和 to() 写回逻辑。

参考资料