在之前的文章中,可能已经对相关概念有所了解。目前SpringBoot集成InfluxDB 2.x时,需要注意InfluxDB 2.x尚未支持SQL查询。对于习惯使用SQL的开发人员来说,直接使用Flux语言可能会存在一定的不适应

目录

一、InfluxDB导入时序数据

二、 SpringBoot集成InfluxDB2.x

1. maven依赖

2. yml配置

3. 配置类

4. 工具类

5. 实体类

6. 模板类

7. 实现类

8. 测试

三、调试过程中踩的坑

1、Flux语句

2、InfluxDB数据转POJO

3、Flux常见报错


        查阅了CSDN上的相关资料,发现大多数文档都是基于Flux实现的。考虑到实际开发需求,决定重新编写相关的工具类来简化操作流程,还是直接操作起来,更容易掌握InfluxDB2.x

 相关文章:

可以参考这两篇文章,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/

官方提供Example :https://github.com/influxdata/influxdb-client-java/blob/9e0ec0be187bdcdab4c03cdb7ded30201e61db6c/client/README.md?plain=1#L647

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 名称

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐