InfluxDB 2.0 Flux 查询语言

FreeGuideOnline 最新 2026-07-11

数据源 -> 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_userusage_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")