InfluxDB 2.0 Flux 查询语言
数据源 -> filter() -> range() -> aggregateWindow() -> yield()
### 表与列结构
Flux 处理的每一张表都有固定的列。时间序列数据被组织成包含 `_time`、`_field`、`_value` 和 `_measurement` 等列的表格。多系列数据会分成多个表,通过**组键(group key)**来区分。
组键是决定表分组的列的集合。例如,按 `_field` 和 `loc` 分组,则每个唯一的 field 和 loc 组合会形成一张独立的表。
---
## 环境准备
在开始编写 Flux 查询前,确保你拥有以下环境之一:
- **InfluxDB 2.0 UI 的 Data Explorer**:直接内建 Flux 脚本编辑器
- **InfluxDB CLI**:使用 `influx query` 命令
- **API 请求**:向 `/api/v2/query` 端点发送 Flux 脚本
本文示例均基于 InfluxDB 2.0 的样本数据(如 `telegraf` 默认 bucket)。你可以在 UI 中加载你自己的数据来测试。
---
## 第一个 Flux 查询
最简单的 Flux 查询包含五个基本部分(五个函数):
```flux
from(bucket: "example-bucket")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "cpu")
|> filter(fn: (r) => r._field == "usage_user")
|> yield(name: "mean")
from():指定数据桶(bucket)range():限定时间范围(支持相对时间如-1h,绝对时间如2022-01-01T00:00:00Z)filter():通过谓词函数筛选数据(r表示当前行记录)yield():输出结果(可省略,但多查询时需要命名)
执行查询:在 Data Explorer 中粘贴后点击 Submit,即可看到表格或图形。
核心查询组件
from() – 指定数据来源
from(bucket: "my-bucket")
还可以从其他数据源读取,例如 CSV 文件:
import "csv"
csv.from(file: "path/to/file.csv")
range() – 掌控时间窗口
绝对时间:
range(start: 2022-01-01T00:00:00Z, stop: 2022-01-02T00:00:00Z)
相对时间:
range(start: -30m) // 过去30分钟
range(start: -1h, stop: -10m) // 1小时前到10分钟前
filter() – 数据筛选
筛选是 Flux 中最常用的操作。谓词函数 fn: (r) => ... 中的 r 代表每一行的列映射,必须返回布尔值。
filter(fn: (r) => r._measurement == "cpu" and r.host == "server01")
使用正则表达式:
filter(fn: (r) => r.host =~ /server[0-9]{2}/)
注意:_value 筛选必须放在 filter() 中,但最好在聚合后使用。
数据转换与整形
map() – 修改行数据
map() 可以创建新列或修改现有列:
from(bucket: "my-bucket")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "cpu")
|> map(fn: (r) => ({ r with
usage_percent: r._value * 100.0,
host_upper: strings.toUpper(v: r.host)
}))
使用 with 关键字保留所有原有列并添加/覆盖指定列。
drop() 和 keep() – 选择列
|> keep(columns: ["_time", "_value", "host"])
|> drop(columns: ["_start", "_stop"])
rename() – 重命名列
|> rename(columns: {_value: "usage", _field: "metric"})
set() – 添加静态列
|> set(key: "source", value: "production")
聚合与窗口操作
aggregateWindow() – 窗口聚合
这是时间序列中最常用的操作之一,用于将数据按时间窗口分组并计算聚合值。
from(bucket: "telegraf")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "cpu" and r._field == "usage_user")
|> aggregateWindow(every: 5m, fn: mean, createEmpty: false)
every:窗口间隔fn:聚合函数(mean,min,max,sum,count,last等)createEmpty:是否保留没有数据的窗口
window() 与手动聚合
window() 会产生按窗口分组的表,配合 mean()/sum() 后需要 duplicate() 或 set 来维护时间列,通常更推荐使用 aggregateWindow 简化操作。
group() 与 ungroup()
group() 修改表的组键,影响后续函数作用范围:
|> group(columns: ["host", "_field"])
ungroup() 移除所有分组,合并为一张表:
|> group()
高级聚合与分析
join() – 表连接
Flux 支持内连接、外连接等,用于合并两个数据流。
典型场景:将 CPU 使用率与内存使用率合并到同一时间线
cpu = from(bucket: "telegraf")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "cpu" and r._field == "usage_user")
|> aggregateWindow(every: 5m, fn: mean)
mem = from(bucket: "telegraf")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "mem" and r._field == "used_percent")
|> aggregateWindow(every: 5m, fn: mean)
join(
tables: {cpu: cpu, mem: mem},
on: ["_time", "host"],
method: "inner"
)
pivot() – 行列转换
将长格式数据转换为宽格式,常用于将多个 field 合并到一行。
from(bucket: "telegraf")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "cpu")
|> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value")
执行后,usage_user 和 usage_system 会变成单独的列。
union() – 垂直合并
union() 将多个结构相同的表纵向拼接。
union(tables: [cpu, mem])
sort() 和 limit()
|> sort(columns: ["_value"], desc: true)
|> limit(n: 10)
变量与自定义函数
定义变量
query_range = -1h
window_interval = 5m
from(bucket: "telegraf")
|> range(start: query_range)
|> aggregateWindow(every: window_interval, fn: mean)
自定义函数
Flux 本身就是由函数组成,定义自己的函数可以大大简化复杂脚本。
getCPUMetrics = (bucket, start, host) =>
from(bucket: bucket)
|> range(start: start)
|> filter(fn: (r) => r._measurement == "cpu" and r.host == host)
getCPUMetrics(bucket: "telegraf", start: -15m, host: "server01")
|> yield()
闭包与高阶函数
可以将 Functions 作为参数传递(如 filter, map 中的 fn),这是 Flux 函数式特性的体现。
处理缺失值与异常
fill() – 填充空值
聚合后,某些窗口可能无数据,可以用 fill() 填充。
|> fill(usePrevious: true) // 用上一个非空值填充
|> fill(value: 0.0) // 用固定值填充
异常值过滤与平滑
使用 movingAverage() 平滑数据:
|> movingAverage(n: 5)
或者用 median() 配合 aggregateWindow 抗离群点。
联合多数据源
Flux 的一个重要优势是可以从多个数据源拉取数据。例如,将 InfluxDB 数据与 CSV 文件合并。
import "csv"
stock = csv.from(url: "https://example.com/stocks.csv")
influx_data = from(bucket: "market")
|> range(start: -1d)
|> filter(fn: (r) => r._measurement == "trades")
join(tables: {influx: influx_data, csv: stock}, on: ["symbol"])
支持的数据源:InfluxDB、CSV、PostgreSQL、MySQL、Prometheus、HTTP 端点等。
实用案例:告警规则与任务
虽然完整的告警和任务定义在 UI 或 API 中,但其核心检测逻辑仍由 Flux 脚本完成。
检测 CPU 超过 80% 并生成告警
from(bucket: "telegraf")
|> range(start: -5m)
|> filter(fn: (r) => r._measurement == "cpu" and r._field == "usage_user")
|> aggregateWindow(every: 1m, fn: mean)
|> map(fn: (r) => ({ r with
level: if r._value > 80.0 then "crit" else "ok"
}))
|> filter(fn: (r) => r.level == "crit")