This commit is contained in:
万佳 2021-09-17 16:35:57 +08:00
commit ea8b9b3d37
11 changed files with 21 additions and 69 deletions

View File

@ -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.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.serialization.StringSerializer;
@ -19,9 +19,6 @@ public class KafkaProducerConfig {
@Value("${spring.kafka.producer.bootstrap_servers:#{null}}") @Value("${spring.kafka.producer.bootstrap_servers:#{null}}")
private String bootstrapServers; private String bootstrapServers;
@Value("${spring.kafka.producer.client_id:#{null}}")
private String clientId;
@Value("${spring.kafka.producer.retries:#{null}}") @Value("${spring.kafka.producer.retries:#{null}}")
private Integer retries; private Integer retries;
@ -30,7 +27,7 @@ public class KafkaProducerConfig {
@Bean @Bean
public KafkaTemplate<String, String> kafkaTemplate() { public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerConfigs(),true); return new KafkaTemplate<>(producerConfigs(), true);
} }
@Bean @Bean
@ -39,7 +36,6 @@ public class KafkaProducerConfig {
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.RETRIES_CONFIG, retries); props.put(ProducerConfig.RETRIES_CONFIG, retries);
props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.CLIENT_ID_CONFIG, clientId);
props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize); props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);

View File

@ -35,31 +35,4 @@ ALTER TABLE gitlink_sys_notification ADD COLUMN (`type` TINYINT(4) NOT NULL DEFA
-- 2021-09-09 新增 source 字段区分消息来源、新增 extra 字段保存额外信息 -- 2021-09-09 新增 source 字段区分消息来源、新增 extra 字段保存额外信息
ALTER TABLE gitlink_sys_notification ADD source varchar(250) NULL COMMENT '消息来源'; ALTER TABLE gitlink_sys_notification ADD source varchar(250) NULL COMMENT '消息来源';
ALTER TABLE gitlink_sys_notification ADD extra TEXT 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;

View File

@ -20,12 +20,6 @@ public class KafkaConsumerConfig {
@Value("${spring.kafka.consumer.bootstrap_servers:#{null}}") @Value("${spring.kafka.consumer.bootstrap_servers:#{null}}")
private String servers; 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}") @Value("${spring.kafka.consumer.auto_offset_reset}")
private String autoOffsetReset; private String autoOffsetReset;
@ -46,9 +40,7 @@ public class KafkaConsumerConfig {
@Bean @Bean
public ConsumerFactory<String, Object> consumerConfigs() { public ConsumerFactory<String, Object> consumerConfigs() {
Map<String, Object> props = new HashMap<>(); Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); 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.MAX_POLL_RECORDS_CONFIG, maxPollRecords);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
props.put(ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG, false); props.put(ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG, false);

View File

@ -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.common.utils.EmailUtils;
import cn.org.gitlink.notification.model.dao.entity.EmailJob; 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.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.EmailJobsService;
import cn.org.gitlink.notification.model.service.notification.EmailSendRecordsService; 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.LogManager;
import org.apache.logging.log4j.Logger; import org.apache.logging.log4j.Logger;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -14,7 +12,6 @@ import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import javax.mail.MessagingException; import javax.mail.MessagingException;
import java.beans.Transient;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;

View File

@ -8,27 +8,23 @@ spring:
kafka: kafka:
producer: producer:
bootstrap_servers: kafka1:9092,kafka2:9092 bootstrap_servers: kafka1:9092,kafka2:9092
client_id: gitlink_producer_01
retries: 5 retries: 5
batch_size: 16384 batch_size: 16384
replication_factor: 1 replication_factor: 1
partitions: 3 partitions: 3
topic_new_email_remind: topic-gitlink-new-email-remind topic_new_email_remind: topic-gitlink-new-email-remind
consumer: consumer:
bootstrap_servers: kafka1:9092,kafka2:9092 bootstrap_servers: kafka1:9092,kafka2:9092
client_id: gitlink_consumer
group_id: group-gitlink-notification group_id: group-gitlink-notification
group_id_email: group-gitlink-email group_id_email: group-gitlink-email
group_id_new_email_remind: group-gitlink-new-email-remind group_id_new_email_remind: group-gitlink-new-email-remind
auto_offset_reset: earliest auto_offset_reset: earliest
max_poll_records: 100 max_poll_records: 100
enable-auto-commit: true
auto-commit-interval: 1S
topic: topic-gitlink-notification topic: topic-gitlink-notification
topic_email: topic-gitlink-email topic_email: topic-gitlink-email
topic_new_email_remind: topic-gitlink-new-email-remind topic_new_email_remind: topic-gitlink-new-email-remind
enable-auto-commit: true
auto-commit-interval: 1S
listener: listener:
concurrency: 3 concurrency: 3
ack-mode: record ack-mode: record

View File

@ -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

View File

@ -1,2 +1 @@
docker-compose -f services.yml down docker-compose -f services.yml down && docker volume prune -f && mvn -f ../pom.xml clean
docker volume prune -f

View File

@ -46,8 +46,8 @@ services:
ZOOKEEPER_TICK_TIME: 2000 ZOOKEEPER_TICK_TIME: 2000
ports: ports:
- ${ZOOKEEPER_LOCAL_PORT}:2181 - ${ZOOKEEPER_LOCAL_PORT}:2181
volumes: # volumes:
- ${DOCKER_DATA_PATH}/zookeeper:/var/lib/zookeeper # - ${DOCKER_DATA_PATH}/zookeeper:/var/lib/zookeeper
networks: networks:
- gitlink_network - gitlink_network
@ -60,8 +60,10 @@ services:
- zookeeper - zookeeper
ports: ports:
- ${KAFKA_01_LOCAL_PORT}:29092 - ${KAFKA_01_LOCAL_PORT}:29092
volumes: # volumes:
- ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_01_NAME}:/var/lib/kafka # - ${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: environment:
KAFKA_BROKER_ID: 1 KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
@ -81,8 +83,10 @@ services:
- zookeeper - zookeeper
ports: ports:
- ${KAFKA_02_LOCAL_PORT}:39092 - ${KAFKA_02_LOCAL_PORT}:39092
volumes: # volumes:
- ${DOCKER_DATA_PATH}/kafka/${KAFKA_CONTAINER_02_NAME}:/var/lib/kafka # - ${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: environment:
KAFKA_BROKER_ID: 2 KAFKA_BROKER_ID: 2
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
@ -105,7 +109,7 @@ services:
environment: environment:
- TZ=Asia/Shanghai - TZ=Asia/Shanghai
volumes: volumes:
- ${DOCKER_DATA_PATH}/logs/:/data/logs/ - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/
depends_on: depends_on:
- kafka1 - kafka1
- kafka2 - kafka2
@ -126,8 +130,7 @@ services:
environment: environment:
- TZ=Asia/Shanghai - TZ=Asia/Shanghai
volumes: volumes:
- ${DOCKER_DATA_PATH}/logs/:/data/logs/ - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/
- TZ=Asia/Shanghai
depends_on: depends_on:
- kafka1 - kafka1
- kafka2 - kafka2
@ -146,7 +149,7 @@ services:
networks: networks:
- gitlink_network - gitlink_network
volumes: volumes:
- ${DOCKER_DATA_PATH}/logs/:/data/logs/ - ${DOCKER_DATA_PATH}/gitlink/:/data/logs/
environment: environment:
- TZ=Asia/Shanghai - TZ=Asia/Shanghai
depends_on: depends_on:

View File

@ -1,2 +1 @@
mvn -f ../pom.xml clean package -DskipTests mvn -f ../pom.xml clean package -DskipTests && docker-compose -f services.yml up --build --force-recreate
docker-compose -f services.yml up --build --force-recreate

View File

@ -21,7 +21,6 @@ spring:
kafka: kafka:
producer: producer:
bootstrap_servers: kafka1:9092,kafka2:9092 bootstrap_servers: kafka1:9092,kafka2:9092
client_id: gitlink_producer_01
retries: 5 retries: 5
batch_size: 16384 batch_size: 16384

View File

@ -8,7 +8,6 @@ spring:
kafka: kafka:
producer: producer:
bootstrap_servers: kafka1:9092,kafka2:9092 bootstrap_servers: kafka1:9092,kafka2:9092
client_id: gitlink_producer_01
retries: 5 retries: 5
batch_size: 16384 batch_size: 16384
replication_factor: 1 replication_factor: 1
@ -28,7 +27,6 @@ spring:
min-idle: 0 min-idle: 0
timeout: 1000 timeout: 1000
datasource: datasource:
driver-class-name: com.mysql.jdbc.Driver 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 url: jdbc:mysql://mysql:3306/gitlink_notification?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&serverTimezone=GMT%2B8&allowMultiQueries=true&useSSL=false