下载并启动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;
    }
}

在这里插入图片描述

Logo

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

更多推荐