diff --git a/wm-data-engine/pom.xml b/wm-data-engine/pom.xml
index d56746cb..a722533e 100644
--- a/wm-data-engine/pom.xml
+++ b/wm-data-engine/pom.xml
@@ -3,15 +3,120 @@
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
4.0.0
- com.waterwm-parent1.0.0-SNAPSHOT
+
+ com.water
+ wm-parent
+ 1.0.0-SNAPSHOT
+
wm-data-engine
+ wm-data-engine
+ 数据汇聚引擎模块
+
- com.waterwm-common
- org.springframework.bootspring-boot-starter-web
- com.alibaba.cloudspring-cloud-starter-alibaba-nacos-discovery
- org.springframework.kafkaspring-kafka
- org.springframework.bootspring-boot-starter-data-redis
- org.postgresqlpostgresql
- net.postgispostgis-jdbc
+
+
+ com.water
+ wm-common
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-websocket
+
+
+
+
+ com.alibaba.cloud
+ spring-cloud-starter-alibaba-nacos-discovery
+
+
+
+
+ org.springframework.kafka
+ spring-kafka
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-data-redis
+
+
+
+
+ org.postgresql
+ postgresql
+
+
+
+
+ net.postgis
+ postgis-jdbc
+
+
+
+
+ com.baomidou
+ mybatis-plus-spring-boot3-starter
+
+
+
+
+ io.minio
+ minio
+
+
+
+
+ cn.hutool
+ hutool-all
+
+
+
+
+ com.github.xiaoymin
+ knife4j-openapi3-jakarta-spring-boot-starter
+
+
+
+
+ com.alibaba
+ easyexcel
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+
+ org.springframework.kafka
+ spring-kafka-test
+ test
+
+
+
+ com.h2database
+ h2
+ test
+
-
\ No newline at end of file
+
+
+
+
+ org.springframework.boot
+ spring-boot-maven-plugin
+
+
+
+
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
new file mode 100644
index 00000000..c56e78bd
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/KafkaConfig.java
@@ -0,0 +1,62 @@
+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.config.ConcurrentKafkaListenerContainerFactory;
+import org.springframework.kafka.core.*;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * Kafka 配置
+ * 用于实时数据流采集和传输
+ */
+@Configuration
+public class KafkaConfig {
+
+ @Value("${spring.kafka.bootstrap-servers:${KAFKA_SERVERS:127.0.0.1}:9092}")
+ 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);
+ }
+
+ @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.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);
+ return factory;
+ }
+}
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/config/MyBatisPlusConfig.java b/wm-data-engine/src/main/java/com/water/data_engine/config/MyBatisPlusConfig.java
new file mode 100644
index 00000000..18cb8dcc
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/MyBatisPlusConfig.java
@@ -0,0 +1,49 @@
+package com.water.data_engine.config;
+
+import com.baomidou.mybatisplus.annotation.DbType;
+import com.baomidou.mybatisplus.core.handlers.MetaObjectHandler;
+import com.baomidou.mybatisplus.extension.plugins.MybatisPlusInterceptor;
+import com.baomidou.mybatisplus.extension.plugins.inner.PaginationInnerInterceptor;
+import org.apache.ibatis.reflection.MetaObject;
+import org.mybatis.spring.annotation.MapperScan;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+
+import java.time.LocalDateTime;
+
+/**
+ * MyBatis-Plus 配置
+ */
+@Configuration
+@MapperScan("com.water.data_engine.mapper")
+public class MyBatisPlusConfig {
+
+ /**
+ * 分页插件
+ */
+ @Bean
+ public MybatisPlusInterceptor mybatisPlusInterceptor() {
+ MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor();
+ interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.POSTGRE_SQL));
+ return interceptor;
+ }
+
+ /**
+ * 自动填充处理器
+ */
+ @Bean
+ public MetaObjectHandler metaObjectHandler() {
+ return new MetaObjectHandler() {
+ @Override
+ public void insertFill(MetaObject metaObject) {
+ this.strictInsertFill(metaObject, "createdAt", LocalDateTime.class, LocalDateTime.now());
+ this.strictInsertFill(metaObject, "updatedAt", LocalDateTime.class, LocalDateTime.now());
+ }
+
+ @Override
+ public void updateFill(MetaObject metaObject) {
+ this.strictUpdateFill(metaObject, "updatedAt", LocalDateTime.class, LocalDateTime.now());
+ }
+ };
+ }
+}
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/config/WebSocketConfig.java b/wm-data-engine/src/main/java/com/water/data_engine/config/WebSocketConfig.java
new file mode 100644
index 00000000..e81045a0
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/config/WebSocketConfig.java
@@ -0,0 +1,32 @@
+package com.water.data_engine.config;
+
+import org.springframework.context.annotation.Configuration;
+import org.springframework.messaging.simp.config.MessageBrokerRegistry;
+import org.springframework.web.socket.config.annotation.EnableWebSocketMessageBroker;
+import org.springframework.web.socket.config.annotation.StompEndpointRegistry;
+import org.springframework.web.socket.config.annotation.WebSocketMessageBrokerConfigurer;
+
+/**
+ * WebSocket 配置
+ * 支持 STOMP 协议,用于实时数据推送
+ */
+@Configuration
+@EnableWebSocketMessageBroker
+public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {
+
+ @Override
+ public void configureMessageBroker(MessageBrokerRegistry registry) {
+ // 客户端订阅前缀: /topic (广播), /queue (点对点)
+ registry.enableSimpleBroker("/topic", "/queue");
+ // 客户端发送消息前缀
+ registry.setApplicationDestinationPrefixes("/app");
+ }
+
+ @Override
+ public void registerStompEndpoints(StompEndpointRegistry registry) {
+ // WebSocket 连接端点
+ registry.addEndpoint("/ws/data-engine")
+ .setAllowedOriginPatterns("*")
+ .withSockJS();
+ }
+}
diff --git a/wm-data-engine/src/main/java/com/water/data_engine/controller/DataCollectController.java b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataCollectController.java
new file mode 100644
index 00000000..83563c3f
--- /dev/null
+++ b/wm-data-engine/src/main/java/com/water/data_engine/controller/DataCollectController.java
@@ -0,0 +1,85 @@
+package com.water.data_engine.controller;
+
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import com.water.common.core.result.R;
+import com.water.data_engine.entity.CollectRecord;
+import com.water.data_engine.entity.CollectTask;
+import com.water.data_engine.service.DataCollectService;
+import io.swagger.v3.oas.annotations.Operation;
+import io.swagger.v3.oas.annotations.tags.Tag;
+import lombok.RequiredArgsConstructor;
+import org.springframework.web.bind.annotation.*;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据采集控制器
+ * DE-01: 实时流(MQTT/Kafka) + 批量采集
+ */
+@Tag(name = "数据采集")
+@RestController
+@RequestMapping("/api/data-engine/collect")
+@RequiredArgsConstructor
+public class DataCollectController {
+
+ private final DataCollectService collectService;
+
+ // ==================== 实时数据采集 ====================
+
+ @Operation(summary = "实时数据接入")
+ @PostMapping("/realtime")
+ public R ingestRealtime(@RequestBody Map request) {
+ String sourceType = (String) request.get("sourceType");
+ String sourceId = (String) request.get("sourceId");
+ @SuppressWarnings("unchecked")
+ Map data = (Map) request.get("data");
+
+ String topic = collectService.ingestRealtime(sourceType, sourceId, data);
+ return R.ok("数据已接入,topic: " + topic);
+ }
+
+ @Operation(summary = "批量数据接入")
+ @PostMapping("/batch")
+ public R batchIngest(@RequestBody List