为什么选择 Spring Boot + InfluxDB?

IoT设备上报的数据是典型的“时序数据”,具有高频、连续、带时间戳的特点。如果使用传统数据库(如MySQL)来处理这类数据,很快就会遇到写入性能瓶颈、存储爆炸等问题。InfluxDB为时序数据场景专门优化,写入吞吐量极高,能以纳秒级精度存储数据,并能按保留策略自动清理过期数据,很适合用于实时监控、日志分析等场景。


📚 InfluxDB 版本与客户端选择

从1.x到2.x,InfluxDB的数据模型发生了根本性的变化:

  • 1.x版本:使用databaseretention policy和InfluxQL查询语言。

  • 2.x版本:引入了organizationbucket的概念,并全面转向Flux查询语言

Spring Boot官方并未提供与2.x版本集成的Starter,因此你需要直接使用官方提供的influxdb-client-java

🚀 项目实战:构建一个智能电表数据监控系统

我们将通过一个完整的示例,展示如何从零开始构建一个接收并查询智能电表上报数据的服务。

1. 初始化项目与添加依赖

首先创建一个Spring Boot项目,并在pom.xml中引入influxdb-client-java依赖:

xml

<dependency>
    <groupId>com.influxdb</groupId>
    <artifactId>influxdb-client-java</artifactId>
    <version>6.12.0</version>
</dependency>
2. 编写配置类 (InfluxDBConfig.java)

我们首先需要创建一个配置类,用于初始化InfluxDB客户端。这种方式可以让我们精细控制客户端的连接参数。我们将连接信息放在application.yml中,然后通过@ConfigurationProperties注入。

2.1 配置文件 (application.yml)

yaml

server:
  port: 8080

influxdb:
  url: http://localhost:8086
  token: 7mJwqFV5H5XQT_KA0FMZbOTyy81La69c4AGIqzyz3CZoOfOBMm1VrPSjv4bqfnkMTzctJ9RjIF8WUZKUM5O10A==
  org: bucketiot
  bucket: bucketiot
  timeout: 10000

smart-meter:
  device-ids:
    - meter-001
    - meter-002
    - meter-003
  data-interval-ms: 5000

2.2 配置类实现

java

@Configuration
public class InfluxDBConfig {

    @Value("${influxdb.url}")
    private String url;

    @Value("${influxdb.token}")
    private String token;

    @Value("${influxdb.org}")
    private String org;

    @Value("${influxdb.bucket}")
    private String bucket;

    @Bean
    public RestTemplate restTemplate() {
        return new RestTemplate();
    }
}

定义一个InfluxDB 数据操作 ,添加查询。

java

@Repository
public class SmartMeterRepository {

    private final RestTemplate restTemplate;
    private final ObjectMapper objectMapper;

    @Value("${influxdb.url}")
    private String url;

    @Value("${influxdb.token}")
    private String token;

    @Value("${influxdb.org}")
    private String org;

    @Value("${influxdb.bucket}")
    private String bucket;

    public SmartMeterRepository(RestTemplate restTemplate) {
        this.restTemplate = restTemplate;
        this.objectMapper = new ObjectMapper();
    }

    public void save(SmartMeterData data) {
        String writeApi = url + "/api/v2/write?org=" + org + "&bucket=" + bucket + "&precision=ms";

        HttpHeaders headers = new HttpHeaders();
        headers.setContentType(MediaType.TEXT_PLAIN);
        headers.set("Authorization", "Token " + token);

        String lineProtocol = buildLineProtocol(data);

        HttpEntity<String> entity = new HttpEntity<>(lineProtocol, headers);
        restTemplate.exchange(writeApi, HttpMethod.POST, entity, String.class);
    }

    private String buildLineProtocol(SmartMeterData data) {
        StringBuilder sb = new StringBuilder();
        sb.append("smart_meter");
        sb.append(",device_id=").append(data.getDeviceId());
        sb.append(",status=").append(data.getStatus());
        sb.append(" voltage=").append(data.getVoltage());
        sb.append(",current=").append(data.getCurrent());
        sb.append(",power=").append(data.getPower());
        sb.append(",energy=").append(data.getEnergy());
        sb.append(",frequency=").append(data.getFrequency());
        sb.append(",power_factor=").append(data.getPowerFactor());
        sb.append(" ").append(data.getTimestamp().toEpochMilli());
        return sb.toString();
    }

    public List<SmartMeterData> query(SmartMeterQuery query) {
        String queryApi = url + "/api/v2/query?org=" + org;

        HttpHeaders headers = new HttpHeaders();
        headers.setContentType(MediaType.APPLICATION_JSON);
        headers.set("Authorization", "Token " + token);

        String fluxQuery = "from(bucket:\"" + bucket + "\") " +
                "|> range(start:" + query.getStartTime().getEpochSecond() + ", stop:" + query.getEndTime().getEpochSecond() + ") " +
                "|> filter(fn:(r)=>r._measurement==\"smart_meter\") " +
                "|> last()";

        if (query.getLimit() != null) {
            fluxQuery += " |> limit(n:" + query.getLimit() + ")";
        }

        String body = "{\"query\":\"" + fluxQuery.replace("\"", "\\\"") + "\",\"type\":\"flux\"}";

        HttpEntity<String> entity = new HttpEntity<>(body, headers);

        try {
            ResponseEntity<String> response = restTemplate.exchange(queryApi, HttpMethod.POST, entity, String.class);
            return parseQueryResult(response.getBody());
        } catch (Exception e) {
            e.printStackTrace();
            return new ArrayList<>();
        }
    }

    private List<SmartMeterData> parseQueryResult(String json) {
        List<SmartMeterData> results = new ArrayList<>();
        if (json == null || json.isEmpty()) {
            return results;
        }

        try {
            JsonNode root = objectMapper.readTree(json);
            JsonNode tables = root.get("response");
            if (tables != null) {
                for (JsonNode table : tables) {
                    JsonNode rows = table.get("series");
                    if (rows != null) {
                        for (JsonNode series : rows) {
                            JsonNode values = series.get("values");
                            if (values != null) {
                                for (JsonNode row : values) {
                                    SmartMeterData data = SmartMeterData.builder()
                                            .deviceId(getStringValue(row, 0))
                                            .timestamp(Instant.ofEpochMilli(getLongValue(row, 1)))
                                            .voltage(getDoubleValue(row, 2))
                                            .current(getDoubleValue(row, 3))
                                            .power(getDoubleValue(row, 4))
                                            .energy(getDoubleValue(row, 5))
                                            .frequency(getDoubleValue(row, 6))
                                            .powerFactor(getDoubleValue(row, 7))
                                            .status(getStringValue(row, 8))
                                            .build();
                                    results.add(data);
                                }
                            }
                        }
                    }
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        }

        return results;
    }

    private String getStringValue(JsonNode row, int index) {
        if (row != null && row.size() > index) {
            JsonNode node = row.get(index);
            return node != null ? node.asText() : null;
        }
        return null;
    }

    private Double getDoubleValue(JsonNode row, int index) {
        if (row != null && row.size() > index) {
            JsonNode node = row.get(index);
            return node != null ? node.asDouble() : null;
        }
        return null;
    }

    private Long getLongValue(JsonNode row, int index) {
        if (row != null && row.size() > index) {
            JsonNode node = row.get(index);
            return node != null ? node.asLong() : null;
        }
        return null;
    }

    public List<SmartMeterData> queryLatestByDevice(String deviceId) {
        SmartMeterQuery query = new SmartMeterQuery();
        query.setDeviceId(deviceId);
        query.setStartTime(Instant.now().minus(java.time.Duration.ofHours(1)));
        query.setEndTime(Instant.now());
        query.setLimit(100);
        return query(query);
    }
4. 编写数据服务 (SmartMeterService.java)

在Service层,  实现写入和查询方法。

java

@Service
public class SmartMeterService {

    private final SmartMeterRepository repository;

    public SmartMeterService(SmartMeterRepository repository) {
        this.repository = repository;
    }

    public void saveData(SmartMeterData data) {
        if (data.getTimestamp() == null) {
            data.setTimestamp(Instant.now());
        }
        repository.save(data);
    }

    public void saveDataBatch(List<SmartMeterData> dataList) {
        for (SmartMeterData data : dataList) {
            if (data.getTimestamp() == null) {
                data.setTimestamp(Instant.now());
            }
            repository.save(data);
        }
    }

    public List<SmartMeterData> queryData(SmartMeterQuery query) {
        if (query.getStartTime() == null) {
            query.setStartTime(Instant.now().minus(java.time.Duration.ofHours(24)));
        }
        if (query.getEndTime() == null) {
            query.setEndTime(Instant.now());
        }
        if (query.getLimit() == null) {
            query.setLimit(100);
        }
        return repository.query(query);
    }

    public List<SmartMeterData> getLatestData(String deviceId) {
        return repository.queryLatestByDevice(deviceId);
    }
5. 编写REST控制器 (SmartMeterController.java)

最后,创建一个REST控制器来对外提供数据上报和查询的API。

java

@RestController
@RequestMapping("/api/smart-meter")
public class SmartMeterController {

    private final SmartMeterService smartMeterService;

    public SmartMeterController(SmartMeterService smartMeterService) {
        this.smartMeterService = smartMeterService;
    }

    @PostMapping("/data")
    public ApiResponse<Void> saveData(@RequestBody SmartMeterData data) {
        smartMeterService.saveData(data);
        return ApiResponse.success(null);
    }

    @PostMapping("/data/batch")
    public ApiResponse<Void> saveDataBatch(@RequestBody List<SmartMeterData> dataList) {
        smartMeterService.saveDataBatch(dataList);
        return ApiResponse.success(null);
    }

    @GetMapping("/data")
    public ApiResponse<List<SmartMeterData>> queryData(
            @RequestParam(required = false) String deviceId,
            @RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE_TIME) Instant startTime,
            @RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE_TIME) Instant endTime,
            @RequestParam(required = false, defaultValue = "100") Integer limit) {

        SmartMeterQuery query = new SmartMeterQuery();
        query.setDeviceId(deviceId);
        query.setStartTime(startTime);
        query.setEndTime(endTime);
        query.setLimit(limit);

        List<SmartMeterData> result = smartMeterService.queryData(query);
        return ApiResponse.success(result);
    }

    @GetMapping("/data/latest/{deviceId}")
    public ApiResponse<List<SmartMeterData>> getLatestData(@PathVariable String deviceId) {
        List<SmartMeterData> result = smartMeterService.getLatestData(deviceId);
        return ApiResponse.success(result);
    }
}

接下来我们启动服务进行测试:

批量上报数据:

可以看到我请求的数据:

我们可以使用InfluxDB2.8 UI界面来查看数据, InfluxDB2.8 UI可以查看官网下载安装:

我们还可以把这些数据做成一个监控看板:

🔍 Flux查询:不止于简单过滤

InfluxDB 2.x 使用 Flux 作为其原生查询和脚本语言。Flux采用函数式管道操作符 (|>) 构建查询,非常灵活且强大。

  • 基础查询from() 指定数据源,range() 定义时间范围,filter() 进行数据过滤。

  • 聚合与转换:Flux提供了丰富的函数,例如使用 mean() 计算平均值,window() 进行时间窗口聚合等。

  • 在UI中调试:在将查询语句放入代码之前,建议先在InfluxDB UI的Data Explorer中编写和调试,确保逻辑正确。


🚨 排查指南:常见问题与解决方案

  • 连接失败或超时:检查URL是否正确、InfluxDB服务是否启动、网络是否通畅。如果使用Docker,检查端口映射是否正确。

  • 401 Unauthorized:Token无效或已过期。请重新生成Token并更新配置。

  • 数据写入失败/无报错:检查orgbucket名称是否正确。WriteApiBlocking的写入操作是同步的,如果失败会抛出异常,可以捕获异常并打印日志以便排查。

  • 批量写入部分成功:关注WriteApi的回调方法,它可以帮助你捕获批量写入中失败的具体数据点。

  • 数据查询不到:检查Flux查询语句中的时间范围(range)是否正确,特别是时区问题。

Logo

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

更多推荐