From cfb4fe6280f2c6f0e636318a75020e8849672a2f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9B=BE=E4=BC=9F?= Date: Sat, 18 Sep 2021 15:11:26 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9kafa=E7=9B=91=E5=90=AC?= =?UTF-8?q?=E5=88=9D=E5=A7=8B=E5=8C=96=E9=85=8D=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../common}/config/KafkaConsumerConfig.java | 28 ++++++++++++------- .../notification/common/utils/KafkaUtil.java | 27 ++++++------------ .../main/resources/application.yml.example | 2 ++ .../main/resources/application.yml.example | 2 ++ 4 files changed, 31 insertions(+), 28 deletions(-) rename {executor/src/main/java/cn/org/gitlink/notification/executor/core => common/src/main/java/cn/org/gitlink/notification/common}/config/KafkaConsumerConfig.java (57%) diff --git a/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaConsumerConfig.java similarity index 57% rename from executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java rename to common/src/main/java/cn/org/gitlink/notification/common/config/KafkaConsumerConfig.java index 6ab75aa..dda9057 100644 --- a/executor/src/main/java/cn/org/gitlink/notification/executor/core/config/KafkaConsumerConfig.java +++ b/common/src/main/java/cn/org/gitlink/notification/common/config/KafkaConsumerConfig.java @@ -1,4 +1,4 @@ -package cn.org.gitlink.notification.executor.core.config; +package cn.org.gitlink.notification.common.config; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.serialization.StringDeserializer; @@ -20,17 +20,21 @@ public class KafkaConsumerConfig { @Value("${spring.kafka.consumer.bootstrap_servers:#{null}}") private String servers; - @Value("${spring.kafka.consumer.auto_offset_reset}") + @Value("${spring.kafka.consumer.auto_offset_reset:#{null}}") private String autoOffsetReset; - @Value("${spring.kafka.consumer.max_poll_records}") + @Value("${spring.kafka.consumer.max_poll_records:#{null}") private String maxPollRecords; - @Value("${spring.kafka.consumer.topic}") + @Value("${spring.kafka.consumer.topic:#{null}}") private String topic; @Bean public ConcurrentKafkaListenerContainerFactory consumerListenerFactory() { + + // 如果当前配置信息内没有kafka consumer则不创建监听工厂 + if (null == servers) return null; + ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory(); factory.setConsumerFactory(consumerConfigs()); factory.setRecordFilterStrategy(record -> record.topic().toLowerCase().equals(this.topic)); @@ -40,12 +44,16 @@ public class KafkaConsumerConfig { @Bean public ConsumerFactory consumerConfigs() { Map props = new HashMap<>(); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); - 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); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + + // 如果当前配置信息内没有kafka consumer,此处解决没有设置配置信息而引起的空指针异常 + if (null != servers) { + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); + 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); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + } return new DefaultKafkaConsumerFactory<>(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 5bf07b8..b4fd419 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 @@ -34,18 +34,15 @@ 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; + @Value("${spring.kafka.producer.topics:''}") + private String topicString; + private AdminClient adminClient; @Autowired @@ -70,20 +67,14 @@ public class KafkaUtil { 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)); - } - + List topics = Arrays.asList(topicString.split(",")); if (!topics.isEmpty()) { - this.createTopic(topics); - + List newTopics = new LinkedList<>(); + topics.forEach(topic -> { + newTopics.add(new NewTopic(topic, partitions, replicationFactor)); + }); + this.createTopic(newTopics); } } diff --git a/executor/src/main/resources/application.yml.example b/executor/src/main/resources/application.yml.example index 3f89c5c..b82de5c 100644 --- a/executor/src/main/resources/application.yml.example +++ b/executor/src/main/resources/application.yml.example @@ -13,6 +13,8 @@ spring: replication_factor: 1 partitions: 3 topic_new_email_remind: topic-gitlink-new-email-remind + # 需要初始化的topic, 添加新的topic后要在这里加上, 用逗号隔开 + topics: ${spring.kafka.producer.topic_new_email_remind} consumer: bootstrap_servers: kafka1:9092,kafka2:9092 group_id: group-gitlink-notification diff --git a/writer/src/main/resources/application.yml.example b/writer/src/main/resources/application.yml.example index 88fc52e..124c9ea 100644 --- a/writer/src/main/resources/application.yml.example +++ b/writer/src/main/resources/application.yml.example @@ -14,6 +14,8 @@ spring: partitions: 3 topic: topic-gitlink-notification topic_email: topic-gitlink-email + # 需要初始化的topic, 添加新的topic后要在这里加上, 用逗号隔开 + topics: ${spring.kafka.producer.topic}, ${spring.kafka.producer.topic_email} redis: database: 0