【Java整合InfluxDB2.0】极简教程,参考官方教程,创建桶,InfluxDBClient,Line Protocol写入数据,查询和封装,配置类,自动获取连接
·
下载并启动influxdb2
-
如下载的为:influxdb2-2.7.6-windows
-
点击 startup.bat,启动的是:influxd.exe
访问并设置密码
- 访问:http://localhost:8086/
设置密码如下:
Username: pf
Password: pf123456
Organization: pf_eit
Bucket: test1 #这里创建的桶不使用
API TOKEN: xxx
导包
- 本教程是 官方的教程,看官方的也一样s
<dependency>
<groupId>com.influxdb</groupId>
<artifactId>influxdb-client-java</artifactId>
<version>6.6.0</version>
</dependency>
创建InfluxDBClient
// You can generate an API token from the "API Tokens Tab" in the UI
String token = System.getenv("INFLUX_TOKEN");
String bucket = "rightcloud";
String org = "pf_eit";
//InfluxDBClient client = InfluxDBClientFactory.create("http://localhost:8086", token.toCharArray());
InfluxDBClient client = InfluxDBClientFactory.create("http://localhost:8086", "pf", "pf123456".toCharArray());
写入
方式1:直接写个字符串。Line Protocol
String data = "mem,host=host1 used_percent=25.43234543";
WriteApiBlocking writeApi = client.getWriteApiBlocking();
writeApi.writeRecord(bucket, org, WritePrecision.NS, data);
方式2:使用key value
Point point = Point
.measurement("mem")
.addTag("host", "host1")
.addField("used_percent", 27.43234543)
.time(Instant.now(), WritePrecision.NS);
WriteApiBlocking writeApi = client.getWriteApiBlocking();
writeApi.writePoint(bucket, org, point);
方式3:使用实体类 未测试
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);
@Measuement(name = "mem") //这个注解,找不到包,就不测了
public static class Mem {
@Column(tag = true)
String host;
@Column
Double used_percent;
@Column(timestamp = true)
Instant time;
}
执行 Flux 查询
String query = "from(bucket: \"rightcloud\") |> range(start: -1h)";
List<FluxTable> tables = client.getQueryApi().query(query, org);
for (FluxTable table : tables) {
for (FluxRecord record : table.getRecords()) {
System.out.println(JSONUtil.toJsonStr(record));
}
}
{
"values": {
"_measurement": "mem",
"_start": 1719824096983,
"result": "_result",
"_field": "used_percent",
"_stop": 1719827696983,
"host": "host1",
"_value": 27.43234543,
"_time": 1719827430644,
"table": 0
},
"table": 0
}
{"values":{"_measurement":"mem","_start":1719824096983,"result":"_result","_field":"used_percent","_stop":1719827696983,"host":"host1","_value":26.43234543,"_time":1719827217425,"table":0},"table":0}
{"values":{"_measurement":"mem","_start":1719824096983,"result":"_result","_field":"used_percent","_stop":1719827696983,"host":"host1","_value":27.43234543,"_time":1719827430644,"table":0},"table":0}
处置客户端 关闭
client.close();
Buckets measurements fields
① 用户名/组织名(pf/pf_eit)
② 数据管理器
③ Buckets 相当于数据库(VibrationSensor)
④ measurements 相当于数据表
⑤ fields 相当于字段
⑥ influxdb自带的chart图表显示工具,可以将选中的数据通过图表显示出来
查看数据
- 登录:http://localhost:8086/
- 选择:BUCKETS——rightcloud(选择这个桶)——点击确定,进入界面
- 选择右侧 Table,选择 mem,点击 SUBMIT,即可看到 表格中数据。
增加配置model
- 增加配置文件
@Data
@Component
@ConfigurationProperties("spring.influx")
public class InfluxdbConfig {
/**
* 请求路径:http://ip:port
*/
private String url;
/**
* 用户名
*/
private String username;
/**
* 密码
*/
private String password;
/**
* 数据库
*/
private String database;
/**
* 保留策略
*/
private String retention;
}
spring:
influx:
url: http://127.0.0.1:8086 #influxdb服务器的地址
username: #用户名
password: #密码
database: rightcloud #指定的数据库
retention_policy: autogen
增加配置类,自动获取连接
@Configuration
public class InfluxDBTemplate {
private final InfluxdbConfig influxdbConfig;
private InfluxDB influxDB;
@Autowired
public InfluxDBTemplate(InfluxdbConfig influxdbConfig) {
this.influxdbConfig = influxdbConfig;
getInfluxDB();
}
/**
* 获取连接
*/
public void getInfluxDB() {
if (influxDB == null) {
//创建连接
influxDB = InfluxDBFactory.connect(influxdbConfig.getUrl(), influxdbConfig.getUsername(), influxdbConfig.getPassword());
// 设置使用数据库,保证库存在
influxDB.setDatabase(influxdbConfig.getDatabase());
// 设置数据库保留策略,保证策略存在
if (ObjectUtils.isEmpty(influxdbConfig.getRetention())) {
influxDB.setRetentionPolicy(influxdbConfig.getRetention());
}
}
}
/**
* 关闭连接
*/
public void close() {
if (influxDB != null) {
influxDB.close();
}
}
查询封装的方法
/**
* select 查询封装
*
* @param queryResult 查询返回结果
* @param clazz 封装对象类型
* @param <T> 泛型
* @return 返回处理回收结果
*/
public <T> List<T> handleQueryResult(QueryResult queryResult, Class<T> clazz) {
// 定义保存结果集合
List<T> lists = new ArrayList<>();
// 获取结果
List<QueryResult.Result> results = queryResult.getResults();
// 遍历结果
results.forEach(result -> {
// 获取 series
List<QueryResult.Series> seriesList = result.getSeries();
// 遍历 series
seriesList.forEach(series -> {
// 获取的所有列
List<String> columns = series.getColumns();
// 获取所有值
List<List<Object>> values = series.getValues();
// 遍历数据 获取结果
for (List<Object> value : values) {
try {
// 根据 clazz 进行封装
T instance = clazz.newInstance();
// 通过 spring 框架提供反射类进行处理
BeanWrapperImpl beanWrapper = new BeanWrapperImpl(instance);
HashMap<String, Object> fields = new HashMap<>();
for (int j = 0; j < columns.size(); j++) {
String column = columns.get(j);
Object val = value.get(j);
if ("time".equals(column)) {
beanWrapper.setPropertyValue("time", Timestamp.from(ZonedDateTime.parse(String.valueOf(val)).toInstant()).getTime());
} else {
// 保存当前列和值到 field map 中
// 注意: 返回结果无须在知道是 tags 还是 fields 认为就是字段和值 可以将所有字段作为 field 进行返回
fields.put(column, val);
}
}
// 通过反射完成 fields 赋值操作
beanWrapper.setPropertyValue("fields", fields);
lists.add(instance);
} catch (InstantiationException | IllegalAccessException e) {
throw new RuntimeException(e);
}
}
});
});
return lists;
}
}

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。
更多推荐

所有评论(0)