From 17d054803583cd870dc15dc16551563ff619c658 Mon Sep 17 00:00:00 2001
From: jpy <812110275@qq.com>
Date: Sat, 27 May 2023 16:12:29 +0800
Subject: [PATCH] Merge remote-tracking branch 'origin/dev' into dev
---
screen-manage/src/main/java/com/moral/api/config/kafka/KafkaConsumerConfig.java | 41 ++++++++++++++++++++++++++++++++++-------
1 files changed, 34 insertions(+), 7 deletions(-)
diff --git a/screen-manage/src/main/java/com/moral/api/config/kafka/KafkaConsumerConfig.java b/screen-manage/src/main/java/com/moral/api/config/kafka/KafkaConsumerConfig.java
index 4d6652a..0a2af8e 100644
--- a/screen-manage/src/main/java/com/moral/api/config/kafka/KafkaConsumerConfig.java
+++ b/screen-manage/src/main/java/com/moral/api/config/kafka/KafkaConsumerConfig.java
@@ -16,8 +16,8 @@
import java.util.HashMap;
import java.util.Map;
-/*@Configuration
-@EnableKafka*/
+@Configuration
+@EnableKafka
public class KafkaConsumerConfig {
@Value("${kafka.consumer.servers}")
private String servers;
@@ -31,21 +31,48 @@
private String autoOffsetReset;
@Value("${kafka.consumer.concurrency}")
private int concurrency;
+ @Value("${kafka.groupId.insert}")
+ private String insertGroupId;
+ @Value("${kafka.groupId.state}")
+ private String stateGroupId;
- @Bean
- public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
+ @Bean("insertListenerContainerFactory")
+ public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> insertListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
- factory.setConsumerFactory(consumerFactory());
+ factory.setConsumerFactory(insertConsumerFactory());
factory.setConcurrency(concurrency);
factory.getContainerProperties().setPollTimeout(1500);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
return factory;
}
- public ConsumerFactory<String, String> consumerFactory() {
- return new DefaultKafkaConsumerFactory<>(consumerConfigs());
+ @Bean("stateListenerContainerFactory")
+ public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> stateListenerContainerFactory() {
+ ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
+ factory.setConsumerFactory(stateConsumerFactory());
+ factory.setConcurrency(concurrency);
+ factory.getContainerProperties().setPollTimeout(1500);
+ factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
+ return factory;
}
+ //kafka������������������������������������
+ public ConsumerFactory<String, String> insertConsumerFactory() {
+ Map<String, Object> map = consumerConfigs();
+ map.put(ConsumerConfig.GROUP_ID_CONFIG, insertGroupId);
+ return new DefaultKafkaConsumerFactory<>(map);
+ }
+
+ //���������������������������������
+ public ConsumerFactory<String, String> stateConsumerFactory() {
+ Map<String, Object> map = consumerConfigs();
+ map.put(ConsumerConfig.GROUP_ID_CONFIG, stateGroupId);
+ return new DefaultKafkaConsumerFactory<>(map);
+ }
+
+ /*
+ * ������������
+ * */
public Map<String, Object> consumerConfigs() {
Map<String, Object> propsMap = new HashMap<>();
propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
--
Gitblit v1.8.0