Flux 查询语法
Flux 是一种面向数据流的查询和处理语言,主要用于 InfluxDB 2.x。它不像 SQL 那样围绕一张二维表编写 SELECT,而是把数据作为“表流”依次传给过滤、分组、变换和聚合函数。
InfluxDB 3 不支持 Flux。新建 InfluxDB 3 系统应使用 SQL 或 InfluxQL;本篇适用于维护 InfluxDB 2.x、已有 Dashboard 和 Task,或阅读旧项目中的 Flux 查询。
最小查询
from(bucket: "iot")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "sensor_data")
|> filter(fn: (r) => r._field == "temperature")执行过程:
读取 iot Bucket
-> 保留最近 1 小时
-> 保留 sensor_data Measurement
-> 保留 temperature Field
-> 返回符合条件的表流|> 称为 Pipe-forward Operator。它把左侧结果作为第一个参数传给右侧函数:
data |> mean()可以理解成:
mean(tables: data)Flux 数据结构
Line Protocol 中的一条 Point:
sensor_data,site=hefei,device_id=sensor01 temperature=26.3,humidity=61.5 1787212800000000000使用 from() 查询后,Field 默认采用“窄表”形式,每个 Field Value 是独立记录:
_time | _measurement | site | device_id | _field | _value |
|---|---|---|---|---|---|
| 08:00 | sensor_data | hefei | sensor01 | temperature | 26.3 |
| 08:00 | sensor_data | hefei | sensor01 | humidity | 61.5 |
核心列:
| 列 | 含义 |
|---|---|
_time | Point 的 Timestamp |
_measurement | Measurement 名称 |
_field | Field Key,例如 temperature |
_value | 对应 Field Value,例如 26.3 |
_start | 当前查询时间范围起点 |
_stop | 当前查询时间范围终点 |
| 其他列 | Tag,例如 site、device_id |
这与 InfluxDB 3 SQL 直接看到 temperature、humidity 两列不同。Flux 中经常先按 _field 过滤;需要把多个 Field 变回普通列时,使用 pivot()。
变量与基本类型
变量
bucketName = "iot"
queryStart = -24h
targetSite = "hefei"
from(bucket: bucketName)
|> range(start: queryStart)
|> filter(fn: (r) => r.site == targetSite)Flux 变量不可重新赋值。变量命名区分大小写。
常用字面量
stringValue = "hefei"
integerValue = 10
floatValue = 26.3
booleanValue = true
durationValue = 5m
timeValue = 2026-08-20T08:00:00Z
arrayValue = ["hefei", "shanghai"]常见 Duration 单位包括 ns、us、ms、s、m、h、d 和 w:
range(start: -15m)
aggregateWindow(every: 1h, fn: mean)Record
Flux 中一行数据称为 Record,通常用 r 表示:
filter(fn: (r) => r.site == "hefei")列名包含特殊字符时使用方括号:
filter(fn: (r) => r["device-id"] == "sensor01")from 与 range
from() 从 Bucket 读取数据:
from(bucket: "iot")从 InfluxDB 读取数据后必须通过 range() 指定时间范围:
// 最近 1 小时
from(bucket: "iot")
|> range(start: -1h)// 固定时间范围,左闭右开
from(bucket: "iot")
|> range(
start: 2026-08-20T08:00:00Z,
stop: 2026-08-20T09:00:00Z,
)// 使用 Dashboard 或 API 提供的时间范围
from(bucket: "iot")
|> range(start: v.timeRangeStart, stop: v.timeRangeStop)range() 不只是普通过滤,它还设置结果的 _start、_stop,这些列通常属于 Group Key。
filter 条件过滤
Measurement、Field 与 Tag
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")也可以合并为一个条件:
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"
)多值条件
from(bucket: "iot")
|> range(start: -24h)
|> filter(fn: (r) =>
r._field == "temperature" or r._field == "humidity"
)使用 contains() 判断数组是否包含某个值:
targetDevices = ["sensor01", "sensor02"]
from(bucket: "iot")
|> range(start: -24h)
|> filter(fn: (r) => contains(value: r.device_id, set: targetDevices))contains() 对大数组可能效率较低。固定的少量值可直接使用 or;大规模动态集合应重新评估 Schema 或查询方案。
正则表达式
from(bucket: "iot")
|> range(start: -24h)
|> filter(fn: (r) => r.device_id =~ /^sensor0[1-3]$/)=~ 表示匹配,!~ 表示不匹配。能用精确比较时优先使用 ==,正则通常需要更多计算。
数值条件
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
只保留需要的列:
from(bucket: "iot")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "sensor_data")
|> keep(columns: ["_time", "_field", "_value", "site", "device_id"])删除指定列:
data
|> drop(columns: ["_start", "_stop"])重命名列:
data
|> rename(columns: {_value: "temperature"})keep() 或 drop() 若移除 Group Key 中的列,会改变输出表的分组结构。
sort、limit 与 tail
按时间升序:
data
|> sort(columns: ["_time"])按时间降序并取前 10 行:
data
|> sort(columns: ["_time"], desc: true)
|> limit(n: 10)取每张输入表的最后 10 行:
data
|> tail(n: 10)这些函数通常对表流中的每张 Table 分别执行,不一定是对整个 Bucket 的全局结果执行。是否需要先 group() 合并表,取决于希望得到“每台设备 10 条”还是“所有设备一共 10 条”。
聚合函数
常用聚合和 Selector:
| 函数 | 作用 |
|---|---|
count() | 统计记录数 |
mean() | 平均值 |
sum() | 求和 |
min()、max() | 最小值、最大值 |
first()、last() | 最早值、最新值 |
median() | 中位数 |
quantile() | 分位数 |
stddev() | 标准差 |
查询每组温度的平均值:
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 内具有相同值。
例如默认可能按以下列分表:
[_start, _stop, _field, _measurement, site, device_id]按站点重新分组:
from(bucket: "iot")
|> range(start: -24h)
|> filter(fn: (r) => r._measurement == "sensor_data")
|> filter(fn: (r) => r._field == "temperature")
|> group(columns: ["site"])
|> mean()此时得到每个站点的平均温度,而不是每台设备的平均温度。
把所有输入合成一个逻辑分组:
data
|> group(columns: [])
|> mean()取消分组会混合不同设备甚至不同 Field,使用前必须确认业务语义。Flux 中许多“结果为什么有很多张表”的问题,都与 Group Key 有关。
aggregateWindow 时间窗口
按 5 分钟计算平均温度:
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 小时取最大值:
data
|> aggregateWindow(every: 1h, fn: max)aggregateWindow() 通常比手工组合 window()、mean() 和 duplicate() 更适合 Dashboard 降采样。
Dashboard 可以使用自动窗口变量:
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 两列:
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() 转为宽表:
data
|> pivot(
rowKey: ["_time"],
columnKey: ["_field"],
valueColumn: "_value",
)
|> keep(columns: ["_time", "site", "device_id", "temperature", "humidity"])结果:
_time | site | device_id | temperature | humidity |
|---|---|---|---|---|
| 08:00 | hefei | sensor01 | 26.3 | 61.5 |
当后续计算需要同时访问多个 Field 时,通常先 pivot():
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 计算新列
将摄氏温度转换为华氏温度:
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 和其他列会丢失。
增加告警级别:
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,并产生一个值。
类型转换
常用转换函数:
data |> toFloat()
data |> toInt()
data |> toString()
data |> toBool()在 map() 中转换某列:
data
|> map(fn: (r) => ({r with retry_count: int(v: r.retry_count)}))转换不能修复已经写入的 Field 类型冲突。更合理的做法是在采集和写入阶段保持类型稳定。
difference、derivative 与 increase
相邻值差
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 可把负差视为异常重置场景,但是否符合计数器语义需要结合数据判断。
单位时间变化率
data
|> derivative(unit: 1m, nonNegative: true)derivative() 根据值差和时间差计算变化率。
计数器增长量
data
|> increase()increase() 常用于可能重置的单调计数器。温度等 Gauge 不适合使用计数器增长函数。
fill 与缺失值
使用固定值替换 null:
data
|> fill(column: "_value", value: 0.0)使用上一条非空值:
data
|> fill(column: "_value", usePrevious: true)若希望为空时间段创建窗口,需要聚合时启用空窗口:
data
|> aggregateWindow(every: 5m, fn: mean, createEmpty: true)
|> fill(column: "_value", usePrevious: true)填充值不是真实采样。设备离线时填 0、沿用上一值或保留空值代表不同业务含义,不能只为了让曲线连续而随意选择。
union 与 join
union 合并相同结构的数据流
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 关联两个数据流
将温度与功率按设备和时间关联:
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 是否包含某列且值非空:
data
|> filter(fn: (r) => exists r.temperature)常与 pivot() 一起使用:
data
|> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value")
|> filter(fn: (r) => exists r.temperature and exists r.humidity)yield 与多个结果
yield() 为输出结果命名:
data
|> mean()
|> yield(name: "mean_temperature")一个脚本可以输出多个结果:
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 与自定义函数
导入标准库包:
import "math"
data
|> map(fn: (r) => ({r with _value: math.round(x: r._value)}))查看 Schema:
import "influxdata/influxdb/schema"
schema.measurements(bucket: "iot")import "influxdata/influxdb/schema"
schema.measurementFieldKeys(
bucket: "iot",
measurement: "sensor_data",
)定义可复用函数:
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:
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 分钟降采样:
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 配置:
| 属性 | 作用 |
|---|---|
name | Task 名称 |
every | 固定间隔执行 |
cron | 使用 Cron 表达式执行,与 every 二选一 |
offset | 延迟执行但保持原计划窗口,用于等待迟到数据 |
offset: 5m 表示整点任务延迟 5 分钟执行,使迟到数据有时间到达;它不是把查询窗口整体向后移动。
创建 Task:
influx task create \
--org example-org \
--file sensor-downsample.fluxTask 需要考虑失败重跑、重复输出、迟到数据和窗口边界。仅使用 range(start: -task.every) 的任务若手动补跑,可能与原定窗口不同;重要场景应结合 option task 提供的运行时间边界设计。
SQL、InfluxQL 与 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:
SELECT "temperature"
FROM "sensor_data"
WHERE time >= now() - 1h
AND "site" = 'hefei';InfluxDB 3 SQL:
SELECT time, site, device_id, temperature
FROM sensor_data
WHERE time >= now() - INTERVAL '1 hour'
AND site = 'hefei';按 5 分钟平均
| 语言 | 时间窗口 |
|---|---|
| Flux | aggregateWindow(every: 5m, fn: mean) |
| InfluxQL | GROUP BY time(5m) + MEAN() |
| InfluxDB 3 SQL | date_bin(INTERVAL '5 minutes', time) + avg() |
常见错误
忘记 range
// 错误:from 后没有时间范围
from(bucket: "iot")
|> filter(fn: (r) => r._measurement == "sensor_data")读取 InfluxDB 时应紧跟 range(),并限制到业务真正需要的范围。
忘记过滤 Field
// 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 全部历史数据再过滤,会制造庞大的中间结果。推荐顺序:
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() 写回逻辑。