diff --git a/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java index 78c737b..d315f31 100644 --- a/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java +++ b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaProducerConfig.java @@ -33,12 +33,17 @@ public class KafkaProducerConfig { @Bean ProducerFactory producerConfigs() { Map props = new HashMap<>(); - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); - props.put(ProducerConfig.RETRIES_CONFIG, retries); - props.put(ProducerConfig.ACKS_CONFIG, "all"); - 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); + + // 如果当前配置信息内没有kafka producer相关配置则不做参数设置,此处解决没有设置配置信息而引起的空指针异常 + if (null != bootstrapServers) { + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); + props.put(ProducerConfig.RETRIES_CONFIG, retries); + props.put(ProducerConfig.ACKS_CONFIG, "all"); + 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); + } + return new DefaultKafkaProducerFactory<>(props); } } diff --git a/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaUtil.java b/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaUtil.java index a43190c..5bf07b8 100644 --- a/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaUtil.java +++ b/common/src/main/java/cn/org/gitlink/notification/common/utils/KafkaUtil.java @@ -14,10 +14,7 @@ import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback; import javax.annotation.PostConstruct; -import java.util.Collection; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; @@ -37,6 +34,18 @@ public class KafkaUtil { @Value("${spring.kafka.producer.bootstrap_servers:#{null}}") private String kafkaServer; + @Value("${spring.kafka.producer.topic:#{null}") + private String topic; + + @Value("${spring.kafka.producer.mail_topic:#{null}") + private String mailTopic; + + @Value("${spring.kafka.producer.partitions:#{null}}") + private Integer partitions; + + @Value("${spring.kafka.producer.replication_factor:#{null}}") + private Short replicationFactor; + private AdminClient adminClient; @Autowired @@ -46,16 +55,38 @@ public class KafkaUtil { * 初始化AdminClient * '@PostConstruct该注解被用来修饰一个非静态的void()方法。 * 被@PostConstruct修饰的方法会在服务器加载Servlet的时候运行,并且只会被服务器执行一次。 - * PostConstruct在构造函数之后执行,init()方法之前执行。 + * PostConstruct在构造函数之后执行,init()方法之前执行。ls */ @PostConstruct private void initAdminClient() { Map props = new HashMap<>(1); + + // 如果当前配置信息内没有kafka producer相关配置则不对adminClient做初始化 if (kafkaServer == null) { return; } + props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServer); adminClient = KafkaAdminClient.create(props); + + // 初始化topics + // TODO: 如果后面新增topic,需要在这里添加 + if (null !=partitions && null != replicationFactor) { + List topics = new LinkedList<>(); + + if (null != topic) { + topics.add(new NewTopic(topic, partitions, replicationFactor)); + } + if (null != mailTopic) { + topics.add(new NewTopic(mailTopic, partitions, replicationFactor)); + } + + if (!topics.isEmpty()) { + this.createTopic(topics); + + } + } + } public void createTopic(Collection newTopics) { diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java index 85c8489..77ba7a9 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java +++ b/executor/src/main/java/cn/org/gitlink/notification/executor/service/jobhandler/EmailJobsListener.java @@ -29,12 +29,6 @@ public class EmailJobsListener { @Value("${spring.kafka.producer.topic_new_email_remind}") private String gitlinkNewEmailRemindTopic; - @Value("${spring.kafka.producer.partitions}") - private Integer partitions; - - @Value("${spring.kafka.producer.replication_factor}") - private Short replicationFactor; - @Autowired private KafkaUtil kafkaUtil; @@ -45,7 +39,6 @@ public class EmailJobsListener { Boolean flag = emailJobsService.sendEmail(newEmailJobVo); //if the message is inserted successfully, send a new email-job message to kafka if (flag){ - kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkNewEmailRemindTopic, partitions, replicationFactor))); kafkaUtil.sendMessage(gitlinkNewEmailRemindTopic, JSONObject.toJSONString(newEmailJobVo)); } } catch (Exception e) { diff --git a/reader/src/main/resources/application.yml.example b/reader/src/main/resources/application.yml.example index ed88882..89cfd92 100644 --- a/reader/src/main/resources/application.yml.example +++ b/reader/src/main/resources/application.yml.example @@ -18,12 +18,6 @@ spring: min-idle: 0 timeout: 1000 - kafka: - producer: - bootstrap_servers: kafka1:9092,kafka2:9092 - retries: 5 - batch_size: 16384 - 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 diff --git a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java index 4a02940..21f801b 100644 --- a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java +++ b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/EmailJobsController.java @@ -10,7 +10,6 @@ import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo; import com.alibaba.fastjson.JSONObject; import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; -import org.apache.kafka.clients.admin.NewTopic; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.springframework.beans.BeanUtils; @@ -22,7 +21,6 @@ import org.springframework.validation.BindingResult; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; -import java.util.Arrays; import java.util.Map; @RestController @@ -35,12 +33,6 @@ public class EmailJobsController { @Value("${spring.kafka.producer.topic_email}") private String gitlinkEmailTopic; - @Value("${spring.kafka.producer.partitions}") - private Integer partitions; - - @Value("${spring.kafka.producer.replication_factor}") - private Short replicationFactor; - @Autowired private KafkaUtil kafkaUtil; @@ -72,7 +64,6 @@ public class EmailJobsController { newEmailJobVo.setPlatform(platform); try { - kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkEmailTopic, partitions, replicationFactor))); kafkaUtil.sendMessage(gitlinkEmailTopic, JSONObject.toJSONString(newEmailJobVo)); return DataPacketUtil.jsonSuccessResult(); } catch (Exception e) { diff --git a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/NotificationController.java b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/NotificationController.java index b818795..478fdb9 100644 --- a/writer/src/main/java/cn/org/gitlink/notification/writer/controller/NotificationController.java +++ b/writer/src/main/java/cn/org/gitlink/notification/writer/controller/NotificationController.java @@ -13,7 +13,6 @@ import cn.org.gitlink.notification.model.service.notification.SysNotificationSer import com.alibaba.fastjson.JSONObject; import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; -import org.apache.kafka.clients.admin.NewTopic; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.springframework.beans.BeanUtils; @@ -25,7 +24,6 @@ import org.springframework.validation.BindingResult; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; -import java.util.Arrays; import java.util.Map; @RestController @@ -36,12 +34,6 @@ public class NotificationController { @Value("${spring.kafka.producer.topic}") private String gitlinkNotificationTopic; - @Value("${spring.kafka.producer.partitions}") - private Integer partitions; - - @Value("${spring.kafka.producer.replication_factor}") - private Short replicationFactor; - @Autowired private KafkaUtil kafkaUtil; @@ -79,7 +71,6 @@ public class NotificationController { newSysNotificationVo.setPlatform(platform); try { - kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkNotificationTopic, partitions, replicationFactor))); kafkaUtil.sendMessage(gitlinkNotificationTopic, JSONObject.toJSONString(newSysNotificationVo)); return DataPacketUtil.jsonSuccessResult(); } catch (Exception e) {