InfluxDB(三)——SpringBoot集成InfluxDB2.x,Flux语言掌握99%
在之前的文章中,可能已经对相关概念有所了解。目前SpringBoot集成InfluxDB 2.x时,需要注意InfluxDB 2.x尚未支持SQL查询。对于习惯使用SQL的开发人员来说,直接使用Flux语言可能会存在一定的不适应
目录
查阅了CSDN上的相关资料,发现大多数文档都是基于Flux实现的。考虑到实际开发需求,决定重新编写相关的工具类来简化操作流程,还是直接操作起来,更容易掌握InfluxDB2.x
相关文章:
- SpringBoot集成InfluxDB 2.x 以及常用方法封装,这篇写的还可以,但是有一些不太熟悉的类,弃之
- InfluxDB查询构建组件,还是觉得有点不够简洁易懂,但有源代码:https://gitee.com/lichenpark/influx-query-wrapper
可以参考这两篇文章,CDSN好多写的什么玩意……
一、InfluxDB导入时序数据
| 版本 | 主要查询语言 | 说明 |
|---|---|---|
| InfluxDB 1.x | InfluxQL | 早期版本使用的 SQL-like 语言。 |
| InfluxDB 2.x | Flux | 当前使用的版本,官方主推 Flux 查询。 |
| InfluxDB 3.x | SQL / InfluxQL | 最新的重构版本,不再支持 Flux,转向原生 SQL |
Flux语句:
data = from(bucket: "example-bucket")
|> range(start: -1h)
|> filter(fn: (r) => r._measurement == "example-measurement" and r._field == "example-field")
看上去还是和SQL区别挺大的,官方文档:https://docs.influxdata.com/influxdb/v2/query-data/flux/#example-data-variable
这个语句的含义是:
从名为 example-bucket 的存储桶中,读取最近 1 小时内,测量名称为 example-measurement 且field字段名为 example-field 的所有时序数据点
- from:从哪一个桶中查询数据
- range:查询最近一小时的数据
- filter: 根据字段、标签或任何其他列值查询数据,类似 SQL 的查询语言中的
SELECT语句和WHERE子句
等同于MySQL语句:
SELECT
_time,
_value
FROM `example-bucket`.`example-measurement`
WHERE _time >= NOW() - INTERVAL 1 HOUR
AND _field = 'example-field';
个人认为没有必要单独去学习Flux语句,在下面实践中就会熟悉Flux语言
本文还是使用官方提供的空气传感器示例数据,因为数据量每天都有
https://docs.influxdata.com/influxdb/v2/reference/sample-data/#air-sensor-sample-data
下载后的数据格式:
从当前表可以看出
- measurement : airSensors
- field:co、humidity、temperature
- tag:sensor_id
三种方式:
1、CSV:https://github.com/influxdata/influxdb2-sample-data/tree/master

2、Cli:
influx write --bucket echola-bucket --url https://influx-testdata.s3.amazonaws.com/air-sensor-data-annotated.csv
3、直接从官方文档下载:(非最新数据)

通过CSV下载只能下载当天的数据,如果保存每天的运行数据,需要创建InfluxDB Task,后面再说……
还是通过InfluxDB Web UI的Load Data——Source——File Upload——Upload a CSV,可参照上一篇文章:InfluxDB(二)——内存原理核心概念以及与MySQL差异通俗解析
二、 SpringBoot集成InfluxDB2.x
踩了好多坑,无语了,利用AI生成了一些代码,还得调半天,还好通了……
只展示部分代码,有些没验证过就不放上来了,一通百通啦,要源代码的可以私我
目录文件大概是:

1. maven依赖
<!-- InfluxDB2 官方客户端 -->
<dependency>
<groupId>com.influxdb</groupId>
<artifactId>influxdb-client-java</artifactId>
<version>5.0.0</version>
</dependency>
<!-- Hutool 工具类 -->
<dependency>
<groupId>cn.hutool</groupId>
<artifactId>hutool-all</artifactId>
<version>5.8.5</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.8.24</version>
</dependency>
2. yml配置
相关参数获取请参照前文
spring:
influxdb:
url: http://localhost:8086
token: your token
org: your org
bucket: your bucket
3. 配置类
package com.echola.influxdblearning.config;
import com.influxdb.client.InfluxDBClient;
import com.influxdb.client.InfluxDBClientFactory;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* @Author: echola
* @Date: 2026/4/21 15:58
* @Description:
*/
@Data
@Configuration
@ConfigurationProperties(prefix = "influx")
public class InfluxDBConfig {
/**
* 连接地址
*/
private String url;
/**
* 认证token
*/
private String token;
/**
* 组织
*/
private String org;
/**
* 数据库
*/
private String bucket;
/**
* 创建 InfluxDB 客户端 Bean
*/
@Bean
public InfluxDBClient influxDBClient() {
return InfluxDBClientFactory.create(url, token.toCharArray(), org, bucket);
}
}
4. 工具类
官方提供的wiki文档:https://deepwiki.com/influxdata/influxdb-client-java/
Flux语句拼接
package com.echola.influxdblearning.utils;
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.DateUtil;
import cn.hutool.core.util.StrUtil;
import cn.hutool.db.sql.Direction;
import com.echola.influxdblearning.enums.FluxEnum;
import com.influxdb.annotations.Measurement;
import java.util.Date;
import java.util.List;
/**
* @Author: echola
* @Date: 2026/4/22 10:16
* @Description: Flux语句拼接
*/
public class FluxUtil {
private final String flux;
public String getFlux() {
return flux;
}
@Override
public String toString() {
return "FluxUtil{" + "flux='" + flux + "'}";
}
private FluxUtil(Builder builder) {
flux = builder.flux.toString();
}
public static class Builder {
private final StringBuilder flux;
private Class<?> measurementClass;
public Builder() {
this.flux = new StringBuilder();
}
public FluxUtil build() {
return new FluxUtil(this);
}
/**
* 设置桶
*/
public Builder bucket(String bucket) {
flux.append("from(bucket: \"").append(bucket).append("\")");
return this;
}
public <T> Builder measurement(Class<T> clazz) {
// 直接从注解拿,不用任何硬编码!
this.measurementClass = clazz;
Measurement measurement = clazz.getAnnotation(Measurement.class);
if (measurement == null) {
throw new IllegalArgumentException("类 " + clazz.getName() + " 未添加 @Measurement 注解");
}
String name = measurement.name();
flux.append(" |> filter(fn: (r) => r._measurement == \"").append(name).append("\")");
return this;
}
/**
* 时间范围查询
*/
public Builder timeRange(Date startTime, Date endTime) {
// Hutool 格式化:2026-04-22T14:30:00Z
String start = DateUtil.format(startTime, DatePattern.UTC_PATTERN);
String end = DateUtil.format(endTime, DatePattern.UTC_PATTERN);
if (StrUtil.isNotBlank(start) && StrUtil.isNotBlank(end)) {
flux.append(" |> range(start: ").append(start).append(", stop: ").append(end).append(")");
} else if (StrUtil.isNotBlank(start)) {
flux.append(" |> range(start: ").append(start).append(")");
} else if (StrUtil.isNotBlank(end)) {
flux.append(" |> range(stop: ").append(end).append(")");
}
return this;
}
/**
* 判断字段是 Tag 还是 Field
*/
private boolean isTag(String fieldName) {
if (measurementClass == null) {
throw new IllegalStateException("请先调用 setMeasurement(Class) 方法");
}
// 遍历所有字段,根据 @Column(name) 匹配,而不是 Java 字段名
for (java.lang.reflect.Field field : measurementClass.getDeclaredFields()) {
com.influxdb.annotations.Column column = field.getAnnotation(com.influxdb.annotations.Column.class);
if (column == null) {
continue;
}
// 获取注解里的数据库列名
String columnName = column.name();
// 如果注解没写 name,则用字段名(兼容默认)
if (columnName.isEmpty()) {
columnName = field.getName();
}
// 匹配传入的 fieldName(如 sensor_id)
if (columnName.equals(fieldName)) {
return column.tag();
}
}
// 遍历完都没找到,抛出异常
throw new IllegalArgumentException("类 " + measurementClass.getName() + " 中不存在字段: " + fieldName);
}
/**
* 条件过滤
*/
public Builder filter(String key, String operator, String value) {
flux.append(" |> filter(fn: (r) => ");
if (isTag(key)) {
// Tag: r.tagKey == "value"
flux.append("r.").append(key).append(" ").append(operator).append(" \"").append(value).append("\"");
} else {
// Field: r._field == "fieldName" and r._value > xxx
flux.append("r._field == \"").append(key)
.append("\" and r._value ").append(operator).append(" ").append(value);
}
flux.append(")");
return this;
}
/**
* 行转列
* 标准写法:按时间分组,将 _field 列展开为多个字段列
*/
public Builder pivot() {
flux.append(" |> pivot(rowKey: [\"_time\"], columnKey: [\"_field\"], valueColumn: \"_value\")");
return this;
}
/**
* 限制返回条数
*/
public Builder limit(long n) {
flux.append(" |> limit(n: ").append(n).append(")");
return this;
}
/**
* 分页
*/
public Builder page(Integer pageNum, Integer pageSize) {
if (pageNum < 1) pageNum = 1;
Integer offset = (pageNum - 1) * pageSize;
flux.append(" |> limit(n: ").append(pageSize).append(", offset: ").append(offset).append(")");
return this;
}
}
}
5. 实体类
映射空气传感器InfluxDB数据
注意此处的Time的类型一定是Instant,InfluxDB2.x不支持String|Date
package com.echola.influxdblearning.entity.po;
import com.influxdb.annotations.Column;
import com.influxdb.annotations.Measurement;
import lombok.Data;
import java.time.Instant;
/**
* @Author: echola
* @Date: 2026/4/21 16:41
* @Description:
*/
@Data
@Measurement(name = "airSensors")
public class AirSensorData {
@Column(name = "sensor_id", tag = true)
private String sensorId;
@Column(name = "temperature")
private Double temperature;
@Column(name = "humidity")
private Double humidity;
@Column(name = "co")
private Double co;
@Column(timestamp = true)
private Instant time;
}
6. 模板类
package com.echola.influxdblearning.template;
import com.echola.influxdblearning.utils.FluxUtil;
import com.influxdb.client.InfluxDBClient;
import com.influxdb.client.QueryApi;
import com.influxdb.client.WriteApi;
import com.influxdb.client.WriteApiBlocking;
import com.influxdb.client.domain.WritePrecision;
import com.influxdb.client.write.Point;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import java.util.List;
/**
* @Author: echola
* @Date: 2026/4/21 18:35
* @Description:
*/
@Component
public class InfluxDBTemplate {
@Autowired
private InfluxDBClient influxDBClient;
public <T> List<T> queryList(FluxUtil.Builder builder, Class<T> clazz) {
String flux = builder.build().getFlux();
return query(flux, clazz);
}
public <T> List<T> query(String flux, Class<T> clazz) {
return queryApi().query(flux, clazz);
}
}
7. 实现类
需求:查询指定时间范围内,传感器TLM0100的数据,按时间倒序,再分页返回
public List<AirSensorDataVO> queryByTimeRange(AirSensorDataDTO dto) {
FluxUtil.Builder builder = new FluxUtil.Builder().bucket(bucket)
.timeRange(dto.getStartTime(), dto.getEndTime())
.measurement(AirSensorData.class)
.filter("sensor_id", "==", dto.getSensorId())
.pivot()
.sort(Direction.DESC)
.page(dto.getPageSize(), dto.getPageNum());
List<AirSensorData> dataList = influxDBTemplate.queryList(builder, AirSensorData.class);
return ConvertUtil.sourceToTarget(dataList, AirSensorDataVO.class);
}
对应的FLux语句
from(bucket: "echola-bucket")
|> range(start: 2026-04-01T08:00:00Z, stop: 2026-04-23T12:00:00Z)
|> filter(fn: (r) => r._measurement == "airSensors")
|> filter(fn: (r) => r.sensor_id == "TLM0100")
|> pivot(rowKey: ["_time"], columnKey: ["_field"], valueColumn: "_value")
|> sort(columns: ["_time"], desc: true)
|> limit(n: 10, offset: 0)
可以在InfluxDB Web UI测试一下Flux语句,数据正常返回

8. 测试
通过PostMan调接口,可以看到数据正确返回啦!!!

三、调试过程中踩的坑
1、Flux语句
关于 Flux,查询语句的顺序非常重要,核心是按顺序串联(管道操作) 各个函数,来处理数据,数据会从一个函数流向另一个函数,就像MySQL也是有顺序的,Select……From……Where……
基本规则与写法
-
固定开端:查询通常以
from()定义数据库桶开始。 -
必须限定时间:
range()函数不可或缺,必须指定查询的时间范围,否则查询不会执行。 -
管道串联:使用
|>符号将一个函数的输出传给下一个函数,像流水线一样处理数据。 -
动态过滤:在
filter()函数中使用fn: (r) => ...的写法来筛选数据行
(1)基本结构
Flux语句的顺序,基本的结构顺序如下:
bucket = "echola-bucket"
start = -1h
stop = now()
from(bucket: bucket)
|> range(start: start, stop: stop)
|> filter(fn: (r) => r._measurement == "airSensors")
|> filter(fn: (r) => r.sensor_id == "TLM001")
|> filter(fn: (r) => r._field == "temperature" or r._field == "humidity" or r._field == "co")
|> aggregateWindow(every: 1m, fn: mean)
|> sort(columns: ["_time"], desc: true)
|> limit(n: 10)
|> yield()
Flux语句严格按照以上顺序执行,不能混淆顺序,日常查询最简固定套路:
from → range → filter → aggregateWindow → sort → limit → yield
虽然Flux是从上往下执行的,但是它不像MySQL一样,先查询所有数据,再进行range和filter。它会做一个关键优化:谓词下推
具体表现是:
range()和filter()这类过滤条件,会被引擎提前执行
引擎直接到存储层,只读取满足时间范围和过滤条件的数据Block,其余数据不会加载内存
所以物理层面,range()和filter()是最先,且同时起作用的,已经进行了一次高效的索引扫描,因此将range、filter等尽快放在前面,它们能被下推到底层存储引擎执行,大幅减少后续处理的数据量
如果你先 filter后 range:
- 数据范围不会被提前裁剪
- 扫描全表
→巨慢 - 结果可能不正确
- 生产环境直接拖垮数据库
那现在再一行一行来看吧,看完Flux掌握99%,哈哈哈哈……
① from(bucket: bucket)(必须)
作用:指定从哪个桶查询,必须写在最前面,没有之一
② |> range(start: start, stop: stop)(必须)
作用:限定时间范围必须紧跟 from!没有它直接报错:unbounded read
时间格式支持:now()、-15s、-1h、-7d、2025-01-01T00:00:00Z
③ |> filter(fn: (r) => r._measurement == "airSensors")(必须)
作用:指定查询哪张表,第一个 filter
④ |> filter(fn: (r) => r.sensor_id == "TLM001")
作用:按 Tag 过滤
⑤ |> filter(fn: (r) => r._field == "temperature" or r._field == "humidity" or r._field == "co")
作用:按照Field过滤,只查询需要的字段,提升速度
如果只需要筛选温度在50~60,那就只用读取temperature这一列
|> filter(fn: (r) => r._field == "temperature")
|> filter(fn: (r) => r._value >= 50 and r._value <= 60)
可以看到Tag是r.tag名称,而field是通过r._field和r._value进行过滤
不要在filter函数里直接做字符串拼接或数学运算,破坏性能优化,仅用于过滤
⑥ |> aggregateWindow(every: 1m, fn: mean)
作用:聚合(1 分钟取一个平均值),用于数据降采样
⑦ |> sort(columns: ["_time"], desc: true)
作用:按时间降序
⑧ |> limit(n: 10)
作用:只返回 10 条,可用于分页。n:每页条数,offset:跳过多少条(从 0 开始)
⑨ |> yield()
作用:表示最终结果返回,单个结果可省略 yield(),多个结果时必须用 yield(name: "结果名") 为每个输出命名
Flux 没有 for 或 if-else,但提供了类似三元运算符的 if 条件表达式(如 a = if true then 1 else 0)
2、InfluxDB数据转POJO
①行转列
由于InfluxDB的数据是列式存储的,不是像MySQL中行式存储,详情可见上篇文章:
InfluxDB 默认存储是长表,直接看真实数据
执行Flux语句:
from(bucket: "echola-bucket")
|> range(start: 2026-04-23T00:00:00Z, stop: 2026-04-24T00:00:00Z)
|> filter(fn: (r) => r._measurement == "airSensors")
|> filter(fn: (r) => r.sensor_id == "TLM0100")
|> sort(columns: ["_time"], desc: true)
|> limit(n: 3, offset: 0)
可以看到数据格式是,数据是:

可以看到是2026-04-23 10:01:51同一个时间点的3条Field(温度、湿度、CO)的数据, 由于不同Field是分开文件存储,并不像MySQL行式存储
_time _field _value sensor_id
--------------------------------------------------------
2026-04-23 10:01:51 temperature 71.38 TLM0100
2026-04-23 10:01:51 humidity 35.19 TLM0100
2026-04-23 10:01:51 co 61.92 TLM0100
再来看看接口的真实效果,还是上面那个接口api/air-sensors/query-by-time-range返回:

返回JSON:
[
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:51"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:41"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:31"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:51"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:41"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:31"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:51"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:41"
},
{
"sensorId": "TLM0100",
"temperature": null,
"humidity": null,
"co": null,
"time": "2026-04-23 10:01:31"
}
]
🔴是不是发现不对了!!!!
理论上应该只返回3行数据,但却返回了9条数据,也缺少了Field字段,跟想象的不一样吧,Web UI是有Field字段的
InfluxDB中筛选后的真实数据是这样的:
time sensor_id _field _value
2026-04-23 10:01:51 TLM0100 temperature 25.5
2026-04-23 10:01:51 TLM0100 humidity 60.0
2026-04-23 10:01:51 TLM0100 co 0.03
2026-04-23 10:01:41 TLM0100 temperature 25.4
2026-04-23 10:01:41 TLM0100 humidity 60.1
2026-04-23 10:01:41 TLM0100 co 0.03
2026-04-23 10:01:31 TLM0100 temperature 25.3
2026-04-23 10:01:31 TLM0100 humidity 60.2
2026-04-23 10:01:31 TLM0100 co 0.03
❓为什么 limit(3) 返回 9 条?
这就跟 一篇文章:InfluxDB(二)——内存原理核心概念以及与MySQL差异通俗解析的InfluxDB的存储结构对应上了,在物理层面上sensor_id=TLM0100是三份完全隔离的物理数据块:
| 存储块标识 | 内容简述 |
|---|---|
| Shard 1 / TSM File | TLM0100 的时间戳列表 + co 压缩值 |
TLM0100 的时间戳列表 + humidity 压缩值 |
|
TLM0100 的时间戳列表 + temperature 压缩值 |
那么3个字段,3个时间点= 3*3 =9条数据,limit 是限制「行」,不是限制「时间点」
由于Field存储在不同的Block,limit依次读取 3 个独立物理文件(temperature、humidity、co),每个文件取前3条,3 + 3 + 3 = 9 条数据
❓为什么Field字段全是 null?
因为数据是“长格式”,属性名(temperature、humidity)在 _field 列里,值在 _value 列里。Java 实体类找的是 temperature 列,当然找不到,映射不上就为 null
理论上应该是——3个时间点的数据3条数据才对
_time temperature humidity co
-----------------------------------------------------
2026-04-23 10:01:51 71.38 35.19 61.92
2026-04-23 10:01:50 71.35 35.20 61.90
2026-04-23 10:01:49 71.32 35.18 61.89
那如何调整成下面这样的格式呢?
✅ 解决办法:必须加 pivot(),它对已加载到内存中的原始数据,按 _time 分组对齐,将多行结构转换为宽表结构
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
含义是:按时间分组,把字段名变成列,把值填入列
最终:
from(bucket: "echola-bucket")
|> range(start: 2026-04-23T00:00:00Z, stop: 2026-04-24T00:00:00Z)
|> filter(fn: (r) => r._measurement == "airSensors")
|> filter(fn: (r) => r.sensor_id == "TLM0100")
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|> sort(columns: ["_time"], desc: true)
|> limit(n: 3, offset: 0)
注意:pivot的顺序,是在filter之后,sort|limit之前
因为pivot会将符合筛序条件的数据,3个文件按照时间合并后,再进行排序和limit。pivot 发生在查询引擎的内存计算层 ,属于数据结构重塑操作,不涉及任何物理存储层的修改或合并
②时间格式是Instant
注意此处的Time的类型一定是Instant,InfluxDB2.x不支持String|Date
2026-04-01T08:00:00Z
Java 对应:Instant.now().toString()
3、Flux常见报错
1. cannot submit unbounded read
原因:没写 range 或顺序错了解决:from 后面紧跟 range
2. undefined record
原因:字段名写错解决:检查 Tag/Field 名称
更多推荐




所有评论(0)