SpringBoot+InfluxDB:IoT数据存储最佳方案
为什么选择 Spring Boot + InfluxDB?
IoT设备上报的数据是典型的“时序数据”,具有高频、连续、带时间戳的特点。如果使用传统数据库(如MySQL)来处理这类数据,很快就会遇到写入性能瓶颈、存储爆炸等问题。InfluxDB为时序数据场景专门优化,写入吞吐量极高,能以纳秒级精度存储数据,并能按保留策略自动清理过期数据,很适合用于实时监控、日志分析等场景。
📚 InfluxDB 版本与客户端选择
从1.x到2.x,InfluxDB的数据模型发生了根本性的变化:
-
1.x版本:使用
database、retention policy和InfluxQL查询语言。 -
2.x版本:引入了
organization、bucket的概念,并全面转向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并更新配置。 -
数据写入失败/无报错:检查
org和bucket名称是否正确。WriteApiBlocking的写入操作是同步的,如果失败会抛出异常,可以捕获异常并打印日志以便排查。 -
批量写入部分成功:关注
WriteApi的回调方法,它可以帮助你捕获批量写入中失败的具体数据点。 -
数据查询不到:检查Flux查询语句中的时间范围(
range)是否正确,特别是时区问题。
更多推荐




所有评论(0)