InfluxDB学习笔记

王梦婵 12 阅读

一、创建连接

1.创建连接

1.InfluxDBClientOptions options = InfluxDBClientOptions.builder()
       .url("http://127.0.0.1:8086")
       .authenticateToken("你的token".toCharArray())
       .org("novaflow")
       .bucket("ems-realtime")
       .writePrecision(WritePrecision.MS)
       .build();
InfluxDBClient client = InfluxDBClientFactory.create(options);

2.关闭数据库连接

if(client != null) client.close();

二、增删改查

1.插入数据

(1)基于行协议写数据
String data = "mem,host=host1 used_percent=23.43234543";

WriteApiBlocking writeApi = client.getWriteApiBlocking();

writeApi.writeRecord(bucket, org, WritePrecision.NS, data);
(2)基于point写数据

数据点(Point)是数据的基本单位,它们被组织在“measurement”中,每个数据点包含时间戳、一个或多个字段(field),以及可选的标签(tag)。

Point point = Point

         .measurement("mem")

         .addTag("host", "host1")

         .addField("used_percent", 23.43234543)

         .time(Instant.now(), WritePrecision.NS);

WriteApiBlocking writeApi = client.getWriteApiBlocking();

writeApi.writePoint(bucket, org, point);
(3)基于pojo写数据
Mem mem = new Mem();

mem.host = "host1";

mem.used_percent = 23.43234543;

mem.time = Instant.now();

WriteApiBlocking writeApi = client.getWriteApiBlocking();

writeApi.writeMeasurement(bucket, org, WritePrecision.NS, mem);

2.查询数据

(1)查询结果映射
①@Measurement对应measurement名称
②@Column:映射字段名
@Measurement(name = "memory")
public class MemoryPoint {

    @Column(name = "time",timestamp = true)
    private Instant time;

    @Column(name = "name",tag=true)
    private String name;

    @Column(name = "free")
    private Long free;

    @Column(name = "used")
    private Long used;

    @Column(name = "buffer")
    private Long buffer;
}
(2)执行查询并映射结果
QueryApi queryApi = client.getQueryApi();
String flux="from(bucket:\"ems-realtime\") |> range(start:-1h) |> filter(fn: (r) => r._measurement == \"memory\")";
List<FluxTable> tables = queryApi.query(flux);
// 2.x POJO映射工具:FluxResultMapper
FluxResultMapper mapper = new FluxResultMapper();
List<MemoryPoint> list = mapper.toPOJOs(tables, MemoryPoint.class);

3.删除数据

(1)InfluxDB 2.x 没有 DROP MEASUREMENT,无法直接重命名数据表;修改名称只能新 measurement 写入,历史数据通过 Flux 迁移。
(2)删除操作可以针对整个measurement、特定标签或时间范围内的数据点进行。
DeleteApi deleteApi = client.getDeleteApi();
Instant start = Instant.EPOCH;
Instant end = Instant.now();
String predicate = "_measurement=\"ems_bms_realtime_monitor\"";
deleteApi.delete(start, end, predicate, "novaflow", "ems-realtime");
(3)命令删除数据,如下删除ems_bms_realtime_monitor数据
influx delete \

--host http://192.168.110.116:8086 \

--token "你的InfluxDB授权Token" \

--org novaflow \

--bucket ems-realtime \

--start 1970-01-01T00:00:00Z \

--stop 2099-12-31T23:59:59Z \

--predicate '_measurement="ems_bms_realtime_monitor"'

三、使用Flux查询InfluxDB

1.三个主要查询函数

(1)From(): 查询 InfluxDB 存储桶中的数据。
(2)Range(): 根据时间限制筛选数据。Flux需要“有界”查询,即查询限制在特定的时间范围内。
(3)Filter(): 根据列值筛选数据。每一行由r表示,每一列由r的一个属性表示。可以应用多个后续筛选器。

Filter ()将每一行读取为名为r的记录。在r记录中,每个键值对表示一列及其值。例如:

r = {

  _time: 2020-01-01T00:00:00Z,

  _measurement: "home",

  room: "Kitchen",

  _field: "temp",

  _value: 21.0,

}

2.查询数据

(1)选择数据点:
from(bucket: "your_bucket")

  |> range(start: -1h)

  |> filter(fn: (r) => r._measurement == "your_measurement")
(2)选择特定字段:
from(bucket: "your_bucket")

  |> range(start: -1h)

  |> filter(fn: (r) => r._measurement == "your_measurement")

  |> filter(fn: (r) => r._field == "your_field")

3.数据处理和计算

(1)计算平均值
from(bucket: "your_bucket")

  |> range(start: -1h)

  |> filter(fn: (r) => r._measurement == "your_measurement")

  |> filter(fn: (r) => r._field == "your_field")

  |> mean()
(2)计算总和
from(bucket: "your_bucket")

  |> range(start: -1h)

  |> filter(fn: (r) => r._measurement == "your_measurement")

  |> filter(fn: (r) => r._field == "your_field")

  |> sum()
(3)计算差异
from(bucket: "your_bucket")

  |> range(start: -1h)

  |> filter(fn: (r) => r._measurement == "your_measurement")

  |> difference()

4.筛选数据

(1)按标签筛选
from(bucket: "your_bucket")

 |> range(start: -1h)

 |> filter(fn: (r) =>
         r._measurement  "your_measurement" 
         and r.your_tag "desired_value"
)
(2)按条件筛选
from(bucket: "your_bucket")

 |> range(start: -1h)

 |> filter(fn: (r) => 
         r._measurement == "your_measurement" 
         and r._value > 10.0
)
(3)按时间范围筛选
from(bucket: "your_bucket")

 |> range(start: -1h, stop: now())

 |> filter(fn: (r) => r._measurement == "your_measurement")

四、总结

本篇笔记主要介绍 InfluxDB2‑x 的 Java 基础操作,它采用组织、存储桶和令牌的鉴权模式。日常开发可通过 Point、POJO、行协议三种方式写入时序数据,依靠 Flux 语句实现数据筛选和聚合统计。该版本并不支持直接删除或是重命名测量表,改名需要迁移历史数据。在 EMS 储能项目里注意复用写入接口、保护访问令牌,规避常见开发问题。