diff --git a/wm-bi/pom.xml b/wm-bi/pom.xml
index dda0db76..2e707090 100644
--- a/wm-bi/pom.xml
+++ b/wm-bi/pom.xml
@@ -13,8 +13,16 @@
cn.dev33sa-token-spring-boot3-starter
org.postgresqlpostgresql
- org.apache.poipoi
- org.apache.poipoi-ooxml
+
+ org.apache.poi
+ poi
+ 5.2.5
+
+
+ org.apache.poi
+ poi-ooxml
+ 5.2.5
+
org.springframework.bootspring-boot-starter-quartz
diff --git a/wm-common/target/maven-status/maven-compiler-plugin/compile/default-compile/createdFiles.lst b/wm-common/target/maven-status/maven-compiler-plugin/compile/default-compile/createdFiles.lst
new file mode 100644
index 00000000..e69de29b
diff --git a/wm-common/target/maven-status/maven-compiler-plugin/compile/default-compile/inputFiles.lst b/wm-common/target/maven-status/maven-compiler-plugin/compile/default-compile/inputFiles.lst
new file mode 100644
index 00000000..9394ed47
--- /dev/null
+++ b/wm-common/target/maven-status/maven-compiler-plugin/compile/default-compile/inputFiles.lst
@@ -0,0 +1,10 @@
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/annotation/DataScope.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/config/SwaggerCommonConfig.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/entity/BaseEntity.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/exception/BusinessException.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/exception/GlobalExceptionHandler.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/result/R.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/storage/MinioService.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/util/ExcelUtils.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/core/util/IdUtils.java
+/root/.openclaw/workspace/water-management-system/wm-common/src/main/java/com/water/common/handler/JsonListTypeHandler.java
diff --git a/wm-data-engine/pom.xml b/wm-data-engine/pom.xml
index a722533e..a071f9fb 100644
--- a/wm-data-engine/pom.xml
+++ b/wm-data-engine/pom.xml
@@ -90,6 +90,13 @@
com.alibaba
easyexcel
+
+
+
+ com.taosdata.jdbc
+ taos-jdbcdriver
+ 3.0.0
+
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/DataEngineApplication.java b/wm-data-engine/src/main/java/com/water/data_engine/DataEngineApplication.java
index 0f6ac0a8..ef4b09c7 100644
--- a/wm-data-engine/src/main/java/com/water/data_engine/DataEngineApplication.java
+++ b/wm-data-engine/src/main/java/com/water/data_engine/DataEngineApplication.java
@@ -1,11 +1,14 @@
package com.water.data_engine;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.context.annotation.Bean;
/**
* 数据引擎应用主类
- * Issue #50: 抄表管理(人工+远传集成)+ 阶梯水价计算
+ * Issue #41: 实时流数据采集(MQTT/Kafka Consumer)
*/
@SpringBootApplication
public class DataEngineApplication {
@@ -13,4 +16,11 @@ public class DataEngineApplication {
public static void main(String[] args) {
SpringApplication.run(DataEngineApplication.class, args);
}
+
+ @Bean
+ public ObjectMapper objectMapper() {
+ ObjectMapper mapper = new ObjectMapper();
+ mapper.registerModule(new JavaTimeModule());
+ return mapper;
+ }
}
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/config/DataEngineInitializer.java b/wm-data-engine/src/main/java/com/water/data_engine/config/DataEngineInitializer.java
new file mode 100644
index 00000000..310fef77
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/DataEngineInitializer.java
@@ -0,0 +1,20 @@
+package com.water.data_engine.config;
+
+import com.water.data_engine.service.TDengineService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.CommandLineRunner;
+import org.springframework.stereotype.Component;
+
+@Component
+public class DataEngineInitializer implements CommandLineRunner {
+
+ @Autowired
+ private TDengineService tdengineService;
+
+ @Override
+ public void run(String... args) throws Exception {
+ // 系统启动时初始化 TDengine 数据库和表
+ tdengineService.initializeDatabase();
+ System.out.println("数据引擎初始化完成");
+ }
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/config/KafkaConfig.java b/wm-data-engine/src/main/java/com/water/data_engine/config/KafkaConfig.java
index c56e78bd..ee0a6a5b 100644
--- a/wm-data-engine/src/main/java/com/water/data_engine/config/KafkaConfig.java
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/KafkaConfig.java
@@ -1,62 +1,45 @@
package com.water.data_engine.config;
import org.apache.kafka.clients.consumer.ConsumerConfig;
-import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
-import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
+import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
-import org.springframework.kafka.core.*;
+import org.springframework.kafka.core.ConsumerFactory;
+import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
+import org.springframework.kafka.support.serializer.JsonDeserializer;
import java.util.HashMap;
import java.util.Map;
-/**
- * Kafka 配置
- * 用于实时数据流采集和传输
- */
+@EnableKafka
@Configuration
public class KafkaConfig {
-
- @Value("${spring.kafka.bootstrap-servers:${KAFKA_SERVERS:127.0.0.1}:9092}")
+
+ @Value("${spring.kafka.bootstrap.servers}")
private String bootstrapServers;
-
- @Bean
- public ProducerFactory producerFactory() {
- Map props = new HashMap<>();
- props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
- props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
- props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
- props.put(ProducerConfig.ACKS_CONFIG, "1");
- props.put(ProducerConfig.RETRIES_CONFIG, 3);
- return new DefaultKafkaProducerFactory<>(props);
- }
-
- @Bean
- public KafkaTemplate kafkaTemplate(ProducerFactory producerFactory) {
- return new KafkaTemplate<>(producerFactory);
- }
-
+
+ @Value("${spring.kafka.consumer.group-id}")
+ private String groupId;
+
@Bean
public ConsumerFactory consumerFactory() {
Map props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
- props.put(ConsumerConfig.GROUP_ID_CONFIG, "wm-data-engine");
+ props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
return new DefaultKafkaConsumerFactory<>(props);
}
-
+
@Bean
- public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory(
- ConsumerFactory consumerFactory) {
- ConcurrentKafkaListenerContainerFactory factory =
- new ConcurrentKafkaListenerContainerFactory<>();
- factory.setConsumerFactory(consumerFactory);
- factory.setConcurrency(3);
+ public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory() {
+ ConcurrentKafkaListenerContainerFactory factory =
+ new ConcurrentKafkaListenerContainerFactory<>();
+ factory.setConsumerFactory(consumerFactory());
return factory;
}
-}
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/config/TDengineConfig.java b/wm-data-engine/src/main/java/com/water/data_engine/config/TDengineConfig.java
new file mode 100644
index 00000000..92b9a8a6
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/TDengineConfig.java
@@ -0,0 +1,60 @@
+package com.water.data_engine.config;
+
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.context.annotation.Configuration;
+
+@Configuration
+@ConfigurationProperties(prefix = "tdengine")
+public class TDengineConfig {
+ private String host;
+ private Integer port = 6030;
+ private String username;
+ private String password;
+ private String database;
+
+ // getters and setters
+ public String getHost() {
+ return host;
+ }
+
+ public void setHost(String host) {
+ this.host = host;
+ }
+
+ public Integer getPort() {
+ return port;
+ }
+
+ public void setPort(Integer port) {
+ this.port = port;
+ }
+
+ public String getUsername() {
+ return username;
+ }
+
+ public void setUsername(String username) {
+ this.username = username;
+ }
+
+ public String getPassword() {
+ return password;
+ }
+
+ public void setPassword(String password) {
+ this.password = password;
+ }
+
+ public String getDatabase() {
+ return database;
+ }
+
+ public void setDatabase(String database) {
+ this.database = database;
+ }
+
+ public String getJdbcUrl() {
+ return String.format("jdbc:TAOS://%s:%d/%s?user=%s&password=%s",
+ host, port, database, username, password);
+ }
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/controller/DataEngineController.java b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataEngineController.java
new file mode 100644
index 00000000..dbcc2e23
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataEngineController.java
@@ -0,0 +1,77 @@
+package com.water.data_engine.controller;
+
+import com.water.data_engine.entity.IotData;
+import com.water.data_engine.service.TDengineService;
+import com.water.data_engine.service.DataCollectService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.web.bind.annotation.*;
+import java.time.LocalDateTime;
+import java.util.HashMap;
+import java.util.Map;
+
+@Slf4j
+@RestController
+@RequestMapping("/api/data-engine")
+public class DataEngineController {
+
+ @Autowired
+ private TDengineService tdengineService;
+
+ @Autowired
+ private DataCollectService dataCollectService;
+
+ @PostMapping("/test-write")
+ public Map testWrite(@RequestBody IotData data) {
+ Map result = new HashMap<>();
+
+ try {
+ // 设置测试数据
+ if (data.getCollectTime() == null) {
+ data.setCollectTime(LocalDateTime.now());
+ }
+ if (data.getStatus() == null) {
+ data.setStatus(1);
+ }
+
+ tdengineService.insertIotData(data);
+
+ result.put("success", true);
+ result.put("message", "测试数据写入成功");
+ result.put("deviceSn", data.getDeviceSn());
+ result.put("collectTime", data.getCollectTime());
+
+ } catch (Exception e) {
+ result.put("success", false);
+ result.put("message", "测试数据写入失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+
+ @GetMapping("/status")
+ public Map getStatus() {
+ Map result = new HashMap<>();
+ result.put("status", "running");
+ result.put("tdengine", "connected");
+ result.put("kafka", "listening");
+ return result;
+ }
+
+ @PostMapping("/initialize")
+ public Map initialize() {
+ Map result = new HashMap<>();
+
+ try {
+ tdengineService.initializeDatabase();
+ result.put("success", true);
+ result.put("message", "TDengine 初始化完成");
+
+ } catch (Exception e) {
+ result.put("success", false);
+ result.put("message", "初始化失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/entity/IotData.java b/wm-data-engine/src/main/java/com/water/data_engine/entity/IotData.java
new file mode 100644
index 00000000..8079cd12
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/entity/IotData.java
@@ -0,0 +1,20 @@
+package com.water.data_engine.entity;
+
+import lombok.Data;
+import java.time.LocalDateTime;
+
+@Data
+public class IotData {
+ private Long id;
+ private String deviceSn;
+ private String deviceType;
+ private Double pressure;
+ private Double flow;
+ private Double temperature;
+ private Double waterLevel;
+ private Double水质指标;
+ private LocalDateTime collectTime;
+ private Integer status;
+ private String location;
+ private String remarks;
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/listener/IotDataKafkaListener.java b/wm-data-engine/src/main/java/com/water/data_engine/listener/IotDataKafkaListener.java
new file mode 100644
index 00000000..1a3901e1
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/listener/IotDataKafkaListener.java
@@ -0,0 +1,48 @@
+package com.water.data_engine.listener;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.water.data_engine.entity.IotData;
+import com.water.data_engine.service.TDengineService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.kafka.annotation.KafkaListener;
+import org.springframework.stereotype.Component;
+import java.time.LocalDateTime;
+
+@Slf4j
+@Component
+public class IotDataKafkaListener {
+
+ @Autowired
+ private TDengineService tdengineService;
+
+ @Autowired
+ private ObjectMapper objectMapper;
+
+ @KafkaListener(topics = "iot-data-topic", groupId = "data-engine-group")
+ public void consumeIotData(String message) {
+ try {
+ log.info("接收到 Kafka 消息: {}", message);
+
+ // 解析 JSON 消息
+ IotData iotData = objectMapper.readValue(message, IotData.class);
+
+ // 设置默认值
+ if (iotData.getCollectTime() == null) {
+ iotData.setCollectTime(LocalDateTime.now());
+ }
+ if (iotData.getStatus() == null) {
+ iotData.setStatus(1); // 默认正常状态
+ }
+
+ // 写入 TDengine
+ tdengineService.insertIotData(iotData);
+
+ log.info("IoT 数据处理完成: 设备={}, 时间={}",
+ iotData.getDeviceSn(), iotData.getCollectTime());
+
+ } catch (Exception e) {
+ log.error("处理 IoT 数据失败: {}", e.getMessage(), e);
+ }
+ }
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/service/TDengineService.java b/wm-data-engine/src/main/java/com/water/data_engine/service/TDengineService.java
new file mode 100644
index 00000000..f8c2268b
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/service/TDengineService.java
@@ -0,0 +1,106 @@
+package com.water.data_engine.service;
+
+import com.water.data_engine.config.TDengineConfig;
+import com.water.data_engine.entity.IotData;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import java.sql.*;
+import java.time.LocalDateTime;
+import java.util.List;
+
+@Service
+public class TDengineService {
+
+ @Autowired
+ private TDengineConfig tdengineConfig;
+
+ public Connection getConnection() throws SQLException {
+ return DriverManager.getConnection(tdengineConfig.getJdbcUrl());
+ }
+
+ public void initializeDatabase() {
+ String createDatabaseSql = String.format("CREATE DATABASE IF NOT EXISTS %s", tdengineConfig.getDatabase());
+ String useDatabaseSql = String.format("USE %s", tdengineConfig.getDatabase());
+ String createTableSql = "CREATE TABLE IF NOT EXISTS iot_data (" +
+ "id BIGINT AUTO_INCREMENT," +
+ "device_sn NCHAR(64) NOT NULL," +
+ "device_type NCHAR(32)," +
+ "pressure DOUBLE," +
+ "flow DOUBLE," +
+ "temperature DOUBLE," +
+ "water_level DOUBLE," +
+ "water_quality_index DOUBLE," +
+ "collect_time TIMESTAMP," +
+ "status INT," +
+ "location NCHAR(128)," +
+ "remarks NCHAR(256)," +
+ "PRIMARY KEY (id, device_sn, collect_time))" +
+ "TAGS (device_type NCHAR(32), location NCHAR(128))";
+
+ try (Connection conn = DriverManager.getConnection(
+ "jdbc:TAOS://" + tdengineConfig.getHost() + ":" + tdengineConfig.getPort() +
+ "?user=" + tdengineConfig.getUsername() + "&password=" + tdengineConfig.getPassword())) {
+
+ Statement stmt = conn.createStatement();
+ stmt.execute(createDatabaseSql);
+ stmt.execute(useDatabaseSql);
+ stmt.execute(createTableSql);
+
+ System.out.println("TDengine 数据库和表初始化完成");
+ } catch (SQLException e) {
+ System.err.println("初始化 TDengine 失败: " + e.getMessage());
+ }
+ }
+
+ public void insertIotData(IotData data) {
+ String sql = String.format("INSERT INTO iot_data VALUES (NULL, '%s', '%s', %.2f, %.2f, %.2f, %.2f, %.2f, '%s', %d, '%s', '%s')",
+ data.getDeviceSn(),
+ data.getDeviceType(),
+ data.getPressure(),
+ data.getFlow(),
+ data.getTemperature(),
+ data.getWaterLevel(),
+ data.get水质指标(),
+ data.getCollectTime().toString(),
+ data.getStatus(),
+ data.getLocation(),
+ data.getRemarks());
+
+ try (Connection conn = getConnection();
+ Statement stmt = conn.createStatement()) {
+ stmt.execute(sql);
+ System.out.println("数据已写入 TDengine: " + data.getDeviceSn());
+ } catch (SQLException e) {
+ System.err.println("写入 TDengine 失败: " + e.getMessage());
+ }
+ }
+
+ public void batchInsertIotData(List dataList) {
+ try (Connection conn = getConnection()) {
+ conn.setAutoCommit(false);
+
+ for (IotData data : dataList) {
+ String sql = String.format("INSERT INTO iot_data VALUES (NULL, '%s', '%s', %.2f, %.2f, %.2f, %.2f, %.2f, '%s', %d, '%s', '%s')",
+ data.getDeviceSn(),
+ data.getDeviceType(),
+ data.getPressure(),
+ data.getFlow(),
+ data.getTemperature(),
+ data.getWaterLevel(),
+ data.get水质指标(),
+ data.getCollectTime().toString(),
+ data.getStatus(),
+ data.getLocation(),
+ data.getRemarks());
+
+ Statement stmt = conn.createStatement();
+ stmt.execute(sql);
+ }
+
+ conn.commit();
+ System.out.println("批量写入 " + dataList.size() + " 条数据到 TDengine");
+ } catch (SQLException e) {
+ System.err.println("批量写入 TDengine 失败: " + e.getMessage());
+ }
+ }
+}
\ No newline at end of file
diff --git a/wm-data-engine/src/main/resources/application.yml b/wm-data-engine/src/main/resources/application.yml
index b9df83f0..0018854d 100644
--- a/wm-data-engine/src/main/resources/application.yml
+++ b/wm-data-engine/src/main/resources/application.yml
@@ -44,6 +44,14 @@ minio:
secret-key: ${MINIO_SECRET_KEY:minioadmin}
bucket: water-management
+# TDengine 配置
+tda:
+ host: ${TDENGINE_HOST:127.0.0.1}
+ port: ${TDENGINE_PORT:6030}
+ username: ${TDENGINE_USER:root}
+ password: ${TDENGINE_PASS:taosdata}
+ database: ${TDENGINE_DB:water_iot}
+
# 日志配置
logging:
level: