diff --git a/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaProducerConfig.java b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java similarity index 85% rename from common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaProducerConfig.java rename to common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java index 869d76e..78c737b 100644 --- a/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaProducerConfig.java +++ b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java @@ -1,4 +1,4 @@ -package cn.org.gitlink.notification.common.utils; +package cn.org.gitlink.notification.common.config; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringSerializer; @@ -19,9 +19,6 @@ public class KafkaProducerConfig { @Value("${spring.kafka.producer.bootstrap_servers:#{null}}") private String bootstrapServers; - @Value("${spring.kafka.producer.client_id:#{null}}") - private String clientId; - @Value("${spring.kafka.producer.retries:#{null}}") private Integer retries; @@ -30,7 +27,7 @@ public class KafkaProducerConfig { @Bean public KafkaTemplate kafkaTemplate() { - return new KafkaTemplate<>(producerConfigs(),true); + return new KafkaTemplate<>(producerConfigs(), true); } @Bean @@ -39,7 +36,6 @@ public class KafkaProducerConfig { props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.RETRIES_CONFIG, retries); props.put(ProducerConfig.ACKS_CONFIG, "all"); - props.put(ProducerConfig.CLIENT_ID_CONFIG, clientId); props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); diff --git a/db/gns-notification.sql b/db/gns-notification.sql index 9f94f9b..1637c77 100644 --- a/db/gns-notification.sql +++ b/db/gns-notification.sql @@ -35,31 +35,4 @@ ALTER TABLE gitlink_sys_notification ADD COLUMN (`type` TINYINT(4) NOT NULL DEFA -- 2021-09-09 新增 source 字段区分消息来源、新增 extra 字段保存额外信息 ALTER TABLE gitlink_sys_notification ADD source varchar(250) NULL COMMENT '消息来源'; -ALTER TABLE gitlink_sys_notification ADD extra TEXT NULL COMMENT '额外信息(备用字段)'; - -DROP TABLE IF EXISTS `gitlink_email_jobs`; -CREATE TABLE `gitlink_email_jobs` ( - `id` int(11) NOT NULL AUTO_INCREMENT, - `sender` int(11) NOT NULL COMMENT '发送者id', - `emails` text NOT NULL COMMENT '收件人全部邮件地址', - `subject` varchar(500) DEFAULT NULL COMMENT '邮件主题', - `content` text NOT NULL COMMENT '邮件内容', - `created_at` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', - `dispatched_at` datetime DEFAULT NULL COMMENT '处理时间', - `dispatched_status` int(11) DEFAULT '-1' COMMENT '发送状态:-1 未处理,1 处理成功,2 处理失败', - PRIMARY KEY (`id`) -) ENGINE=InnoDB DEFAULT CHARSET=utf8mb3; - - -DROP TABLE IF EXISTS `gitlink_email_send_records`; -CREATE TABLE `gitlink_email_send_records` ( - `id` int(11) NOT NULL AUTO_INCREMENT, - `email` varchar(500) NOT NULL DEFAULT '' COMMENT '收件人邮件地址', - `job_id` int(11) NOT NULL COMMENT '邮件详情id', - `created_at` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', - `sent_at` datetime DEFAULT NULL COMMENT '发送时间', - `status` int(11) DEFAULT '-1' COMMENT '发送状态:-1 未发送,1 发送成功,2 发送失败', - PRIMARY KEY (`id`), - KEY `index_on_email_and_status` (`email`,`status`), - KEY `index_on_status` (`status`) -) ENGINE=InnoDB DEFAULT CHARSET=utf8mb3; \ No newline at end of file +ALTER TABLE gitlink_sys_notification ADD extra TEXT NULL COMMENT '额外信息(备用字段)'; \ No newline at end of file diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java b/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java index 01db58b..6ab75aa 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java @@ -20,12 +20,6 @@ public class KafkaConsumerConfig { @Value("${spring.kafka.consumer.bootstrap_servers:#{null}}") private String servers; - @Value("${spring.kafka.consumer.group_id}") - private String groupId; - - @Value("${spring.kafka.consumer.client_id}") - private String clientId; - @Value("${spring.kafka.consumer.auto_offset_reset}") private String autoOffsetReset; @@ -46,9 +40,7 @@ public class KafkaConsumerConfig { @Bean public ConsumerFactory consumerConfigs() { Map props = new HashMap<>(); - props.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); - props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset); props.put(ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG, false); diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java b/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java index 63f2747..f6082a7 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/service/email/EmailService.java @@ -3,10 +3,8 @@ package cn.org.gitlink.notification.executor.service.email; import cn.org.gitlink.notification.common.utils.EmailUtils; import cn.org.gitlink.notification.model.dao.entity.EmailJob; import cn.org.gitlink.notification.model.dao.entity.EmailSendRecord; -import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo; import cn.org.gitlink.notification.model.service.notification.EmailJobsService; import cn.org.gitlink.notification.model.service.notification.EmailSendRecordsService; -import com.alibaba.fastjson.JSONObject; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.springframework.beans.factory.annotation.Autowired; @@ -14,7 +12,6 @@ import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Transactional; import javax.mail.MessagingException; -import java.beans.Transient; import java.util.ArrayList; import java.util.Date; import java.util.List; diff --git a/executor/src/main/resources/application.yml.example b/executor/src/main/resources/application.yml.example index 6dd55fd..3f89c5c 100644 --- a/executor/src/main/resources/application.yml.example +++ b/executor/src/main/resources/application.yml.example @@ -8,27 +8,23 @@ spring: kafka: producer: bootstrap_servers: kafka1:9092,kafka2:9092 - client_id: gitlink_producer_01 retries: 5 batch_size: 16384 replication_factor: 1 partitions: 3 topic_new_email_remind: topic-gitlink-new-email-remind - consumer: bootstrap_servers: kafka1:9092,kafka2:9092 - client_id: gitlink_consumer group_id: group-gitlink-notification group_id_email: group-gitlink-email group_id_new_email_remind: group-gitlink-new-email-remind auto_offset_reset: earliest max_poll_records: 100 + enable-auto-commit: true + auto-commit-interval: 1S topic: topic-gitlink-notification topic_email: topic-gitlink-email topic_new_email_remind: topic-gitlink-new-email-remind - enable-auto-commit: true - auto-commit-interval: 1S - listener: concurrency: 3 ack-mode: record diff --git a/middleware/end_docker_compose.bat b/middleware/end_docker_compose.bat index d8924a7..b06c53d 100644 --- a/middleware/end_docker_compose.bat +++ b/middleware/end_docker_compose.bat @@ -1 +1 @@ -docker-compose -f services.yml down && docker volume prune -f +docker-compose -f services.yml down && docker volume prune -f && mvn -f ../pom.xml clean diff --git a/middleware/end_docker_compose.sh b/middleware/end_docker_compose.sh index ae55462..51e7dcd 100755 --- a/middleware/end_docker_compose.sh +++ b/middleware/end_docker_compose.sh @@ -1,2 +1 @@ -docker-compose -f services.yml down -docker volume prune -f +docker-compose -f services.yml down && docker volume prune -f && mvn -f ../pom.xml clean diff --git a/middleware/services.yml b/middleware/services.yml index e7b81b8..bb1bb8f 100644 --- a/middleware/services.yml +++ b/middleware/services.yml @@ -46,8 +46,8 @@ services: ZOOKEEPER_TICK_TIME: 2000 ports: - ${ZOOKEEPER_LOCAL_PORT}:2181 - volumes: - - ${DOCKER_DATA_PATH}/zookeeper:/var/lib/zookeeper +# volumes: +# - ${DOCKER_DATA_PATH}/zookeeper:/var/lib/zookeeper networks: - gitlink_network @@ -60,8 +60,10 @@ services: - zookeeper ports: - ${KAFKA_01_LOCAL_PORT}:29092 - volumes: - - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_01_NAME}:/var/lib/kafka +# volumes: +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_01_NAME}/lib:/var/lib/kafka +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_01_NAME}/logs:/var/logs/kafka +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_01_NAME}/conf:/etc/kafka environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 @@ -81,8 +83,10 @@ services: - zookeeper ports: - ${KAFKA_02_LOCAL_PORT}:39092 - volumes: - - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_02_NAME}:/var/lib/kafka +# volumes: +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_02_NAME}/lib:/var/lib/kafka +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_02_NAME}/logs:/var/logs/kafka +# - ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_02_NAME}/conf:/etc/kafka environment: KAFKA_BROKER_ID: 2 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 @@ -105,7 +109,7 @@ services: environment: - TZ=Asia/Shanghai volumes: - - ${DOCKER_DATA_PATH}/logs/:/data/logs/ + - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/ depends_on: - kafka1 - kafka2 @@ -126,8 +130,7 @@ services: environment: - TZ=Asia/Shanghai volumes: - - ${DOCKER_DATA_PATH}/logs/:/data/logs/ - - TZ=Asia/Shanghai + - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/ depends_on: - kafka1 - kafka2 @@ -146,7 +149,7 @@ services: networks: - gitlink_network volumes: - - ${DOCKER_DATA_PATH}/logs/:/data/logs/ + - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/ environment: - TZ=Asia/Shanghai depends_on: diff --git a/middleware/start_docker_compose.sh b/middleware/start_docker_compose.sh index 1d22359..f5bfb1b 100755 --- a/middleware/start_docker_compose.sh +++ b/middleware/start_docker_compose.sh @@ -1,2 +1 @@ -mvn -f ../pom.xml clean package -DskipTests -docker-compose -f services.yml up --build --force-recreate \ No newline at end of file +mvn -f ../pom.xml clean package -DskipTests && docker-compose -f services.yml up --build --force-recreate \ No newline at end of file diff --git a/reader/src/main/resources/application.yml.example b/reader/src/main/resources/application.yml.example index 5e49bde..ed88882 100644 --- a/reader/src/main/resources/application.yml.example +++ b/reader/src/main/resources/application.yml.example @@ -21,7 +21,6 @@ spring: kafka: producer: bootstrap_servers: kafka1:9092,kafka2:9092 - client_id: gitlink_producer_01 retries: 5 batch_size: 16384 diff --git a/writer/src/main/resources/application.yml.example b/writer/src/main/resources/application.yml.example index ec4217b..88fc52e 100644 --- a/writer/src/main/resources/application.yml.example +++ b/writer/src/main/resources/application.yml.example @@ -8,7 +8,6 @@ spring: kafka: producer: bootstrap_servers: kafka1:9092,kafka2:9092 - client_id: gitlink_producer_01 retries: 5 batch_size: 16384 replication_factor: 1 @@ -28,7 +27,6 @@ spring: min-idle: 0 timeout: 1000 - datasource: driver-class-name: com.mysql.jdbc.Driver url: jdbc:mysql://mysql:3306/gitlink_notification?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&serverTimezone=GMT%2B8&allowMultiQueries=true&useSSL=false