springboot集成kafka消费手动启动停止操作

 更新时间:2022年09月07日 11:28:29   作者:zengliangxi  
这篇文章主要介绍了springboot集成kafka消费手动启动停止操作,本文给大家介绍项目场景及解决分析,结合实例代码给大家介绍的非常详细,需要的朋友可以参考下

项目场景:

在月结,或者某些时候,我们需要停掉kafka所有的消费端,让其暂时停止消费,而后等月结完成,再从新对消费监听恢复,进行消费,此动作不需要重启服务

解决分析

KafkaListenerEndpointRegistry这是kafka与spring集成的监听注册bean,可以通过它获取监听容器对象,然后对监听容器对象实行启动,暂停,恢复等操作

/**
 * kafka服务操作类
 * @author liangxi.zeng
 */
@Service
@Slf4j
public class KafkaService {

    @Autowired
    private KafkaListenerEndpointRegistry registry;

    /**
     * 开启消费
     * @param listenerId
     */
    public void start(String listenerId) {
        MessageListenerContainer messageListenerContainer = registry
                .getListenerContainer(listenerId);
        if(Objects.nonNull(messageListenerContainer)) {
            if(!messageListenerContainer.isRunning()) {
                messageListenerContainer.start();
            } else {
                if(messageListenerContainer.isContainerPaused()) {
                    log.info("listenerId:{},恢复",listenerId);
                    messageListenerContainer.resume();
                }
            }
        }
    }

    /**
     * 停止消费
     * @param listenerId
     */
    public void pause(String listenerId) {
        MessageListenerContainer messageListenerContainer = registry
                .getListenerContainer(listenerId);
        if(Objects.nonNull(messageListenerContainer) && !messageListenerContainer.isContainerPaused()) {
            log.info("listenerId:{},暂停",listenerId);
            messageListenerContainer.pause();
        }
    }
}

kafka启动,停止,恢复触发场景

1.通过定时任务自动触发,通过@Scheduled,在某个时间点暂停kafka某个监听的消费,也可以在某个时间点恢复kafka某个监听的消费

/**
 * 卡夫卡配置类
 * @author liangxi.zeng
 */
@Configuration
@EnableScheduling
public class KafkaConfigure {

    @Autowired
    private KafkaService kafkaService;

    @Autowired
    private KafkaConfigParam kafkaConfigParam;

   @Scheduled(cron = "0/10 * * * * ?")
    public void startListener() {
        List<String> topics = kafkaConfigParam.getStartTopics();
        System.out.println("开启。。。"+topics);
        Optional.ofNullable(topics).orElse(new ArrayList<>()).forEach(topic -> {
            kafkaService.start(topic);
        });

    }

    @Scheduled(cron = "0/10 * * * * ?")
    public void pauseListener() {
        List<String> topics = kafkaConfigParam.getPauseTopics();
        System.out.println("暂停。。。"+topics);
        Optional.ofNullable(topics).orElse(new ArrayList<>()).forEach(topic -> {
            kafkaService.pause(topic);
        });

    }
}

2.通过访问接口手动触发kafka消费的启动,暂停,恢复

@RequestMapping("/start/{kafkaId}")
    public String start(@PathVariable String kafkaId) {
        if(!registry.getListenerContainer(kafkaId).isRunning()) {
            registry.getListenerContainer(kafkaId).start();
        } else {
            registry.getListenerContainer(kafkaId).resume();
        }

        return "ok";
    }

    @RequestMapping("/pause/{kafkaId}")
    public String pause(@PathVariable String kafkaId) {
        registry.getListenerContainer(kafkaId).pause();
        return "ok";
    }

3.监听nacos配置文件,完成动态的启停操作

/**
 * nacos配置变更监听
 * @author liangxi.zeng
 */
@Component
@Slf4j
public class NacosConfigListener {

    @Autowired
    private NacosConfigManager nacosConfigManager;

    @Autowired
    private KafkaService kafkaService;

    @Autowired
    private KafkaStartPauseParam kafkaStartPauseParam;

    /**
     * 分隔符
     */
    private static final String SPLIT=",";

    private static final String GROUP = "DEFAULT_GROUP";


    /**
     * nacos 配置文件监听
     * @throws NacosException
     */
    @PostConstruct
    private void reloadConfig() throws NacosException {
        nacosConfigManager.getConfigService().addListener(kafkaStartPauseParam.getDataId(), GROUP, new AbstractConfigChangeListener() {
            @Override
            public void receiveConfigChange(final ConfigChangeEvent event) {
                ConfigChangeItem pauseListeners = event.getChangeItem(KafkaStartPauseParam.PREFIX+".pause-listener");
                ConfigChangeItem startListeners = event.getChangeItem(KafkaStartPauseParam.PREFIX+".start-listener");
                if(Objects.nonNull(pauseListeners)) {
                    pause(pauseListeners);
                }

                if(Objects.nonNull(startListeners)) {
                    start(startListeners);
                }
            }
        });
    }

    /**
     * 暂停消费
     * @param pauseListeners
     */
    private void pause(ConfigChangeItem pauseListeners) {
        String pauseValue = pauseListeners.getNewValue();
        log.info("暂停listener:{}",pauseValue);
        if(!StringUtils.isEmpty(pauseValue)) {
            String[] pauseListenerIds = pauseValue.split(SPLIT);
            for(String pauseListenerId:pauseListenerIds) {
                kafkaService.pause(pauseListenerId);
            }
        }
    }

    /**
     * 恢复消费
     * @param startListeners
     */
    private void start(ConfigChangeItem startListeners) {
        String startValue = startListeners.getNewValue();
        log.info("启动listener:{}",startValue);
        if(!StringUtils.isEmpty(startValue)) {
            String[] startListenerIds = startValue.split(SPLIT);
            for(String startListenerId:startListenerIds) {
                kafkaService.start(startListenerId);
            }
        }
    }

}

配置类

/**
 * kafka配置参数
 * @author liangxi.zeng
 */
@ConfigurationProperties(prefix = KafkaStartPauseParam.PREFIX)
@Data
@Component
@RefreshScope
public class KafkaStartPauseParam {

    public static final String PREFIX = "tcl.kafka";

    private String pauseListener;

    private String startListener;

    private String dataId;
}

源码分析

1.springboot集成kafka,集成配置类org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration

2.@Import({KafkaAnnotationDrivenConfiguration.class})

@Configuration(
        proxyBeanMethods = false
    )
    @EnableKafka
    @ConditionalOnMissingBean(
        name = {"org.springframework.kafka.config.internalKafkaListenerAnnotationProcessor"}
    )
    static class EnableKafkaConfiguration {
        EnableKafkaConfiguration() {
        }
    }
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Import(KafkaListenerConfigurationSelector.class)
public @interface EnableKafka {
}
@Order
public class KafkaListenerConfigurationSelector implements DeferredImportSelector {

	@Override
	public String[] selectImports(AnnotationMetadata importingClassMetadata) {
		return new String[] { KafkaBootstrapConfiguration.class.getName() };
	}

}
public class KafkaBootstrapConfiguration implements ImportBeanDefinitionRegistrar {

	@Override
	public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, BeanDefinitionRegistry registry) {
		if (!registry.containsBeanDefinition(
				KafkaListenerConfigUtils.KAFKA_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)) {

			registry.registerBeanDefinition(KafkaListenerConfigUtils.KAFKA_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME,
					new RootBeanDefinition(KafkaListenerAnnotationBeanPostProcessor.class));
		}

		if (!registry.containsBeanDefinition(KafkaListenerConfigUtils.KAFKA_LISTENER_ENDPOINT_REGISTRY_BEAN_NAME)) {
			registry.registerBeanDefinition(KafkaListenerConfigUtils.KAFKA_LISTENER_ENDPOINT_REGISTRY_BEAN_NAME,
					new RootBeanDefinition(KafkaListenerEndpointRegistry.class));
		}
	}

}

3.KafkaListenerAnnotationBeanPostProcessor这个类,就是消费监听的解析类,在此类中,将监听的方法放入了监听容器MessageListenerContainer

4.监听容器中有ListenerConsumer监听消费者的属性,此内部内实现了SchedulingAwareRunnable接口,此接口继承了Runnable接口,完成了定时异步消费等操作

@Override
		public void run() {
			
			while (isRunning()) {
				try {
					pollAndInvoke();
				}
			}
			wrapUp();
		}

protected void pollAndInvoke() {
			if (!this.autoCommit && !this.isRecordAck) {
				processCommits();
			}
			idleBetweenPollIfNecessary();
			if (this.seeks.size() > 0) {
				processSeeks();
			}
			pauseConsumerIfNecessary();
			this.lastPoll = System.currentTimeMillis();
			this.polling.set(true);
			ConsumerRecords<K, V> records = doPoll();
			if (!this.polling.compareAndSet(true, false) && records != null) {
				/*
				 * There is a small race condition where wakeIfNecessary was called between
				 * exiting the poll and before we reset the boolean.
				 */
				if (records.count() > 0) {
					this.logger.debug(() -> "Discarding polled records, container stopped: " + records.count());
				}
				return;
			}
			resumeConsumerIfNeccessary();
			debugRecords(records);
			if (records != null && records.count() > 0) {
				if (this.containerProperties.getIdleEventInterval() != null) {
					this.lastReceive = System.currentTimeMillis();
				}
				invokeListener(records);
			}
			else {
				checkIdle();
			}
		}

遗留问题

在对kafka消费监听启停的过程中,发现当暂停消费的时候,对于存量的topic还是会消费完,不会立即停止,只是对于新产生的topic不会再消费了

到此这篇关于springboot集成kafka消费手动启动停止的文章就介绍到这了,更多相关springboot集成kafka内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!

相关文章

  • java实现上传文件类型检测过程解析

    java实现上传文件类型检测过程解析

    这篇文章主要介绍了java实现上传文件类型检测过程解析,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下
    2019-12-12
  • Springboot中使用拦截器、过滤器、监听器的流程分析

    Springboot中使用拦截器、过滤器、监听器的流程分析

    Javaweb三大组件:servlet、Filter(过滤器)、 Listener(监听器),这篇文章主要介绍了Springboot中使用拦截器、过滤器、监听器的流程分析,感兴趣的朋友跟随小编一起看看吧
    2023-12-12
  • Java并发编程ThreadLocalRandom类详解

    Java并发编程ThreadLocalRandom类详解

    这篇文章主要介绍了Java并发编程ThreadLocalRandom类详解,通过提出问题为什么需要ThreadLocalRandom展开详情,感兴趣的朋友可以参考一下
    2022-06-06
  • Java中Boolean与字符串或者数字1和0的转换实例

    Java中Boolean与字符串或者数字1和0的转换实例

    下面小编就为大家带来一篇Java中Boolean与字符串或者数字1和0的转换实例。小编觉得挺不错的,现在就分享给大家,也给大家做个参考。一起跟随小编过来看看吧
    2017-07-07
  • SpringBoot 自动装配的原理详解分析

    SpringBoot 自动装配的原理详解分析

    这篇文章主要介绍了SpringBoot 自动装配的原理详解分析,文章通过通过一个案例来看一下自动装配的效果展开详情,感兴趣的小伙伴可以参考一下
    2022-08-08
  • IntelliJ IDEA2020.1 Mac maven sdk 全局配置

    IntelliJ IDEA2020.1 Mac maven sdk 全局配置

    这篇文章主要介绍了IntelliJ IDEA2020.1 Mac maven sdk 全局配置,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2020-06-06
  • idea创建maven父子工程导致子工程无法导入父工程依赖

    idea创建maven父子工程导致子工程无法导入父工程依赖

    创建maven父子工程时遇到一个问题,本文主要介绍了idea创建maven父子工程导致子工程无法导入父工程依赖,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2022-04-04
  • 深入理解Java对象复制

    深入理解Java对象复制

    使用任何已有的工具,都没有直接使用 get set 方式进行,对象转换的速度快,虽然get set 方式代码对一些比较麻烦,但是效率要高一些的,推荐使用 MapStruct 方式.,需要的朋友可以参考下
    2021-05-05
  • Java使用Spring AI的10个实用技巧分享

    Java使用Spring AI的10个实用技巧分享

    在当今的软件开发领域,人工智能的集成已经成为提升应用功能和用户体验的重要手段,对于 Java 开发者而言,Spring AI 提供了一套强大且便捷的工具,本文将深入探讨 Java 使用 Spring AI 的 10 个实用技巧,感兴趣的小伙伴跟着小编一起来看看吧
    2025-06-06
  • Jmeter跨线程组传值调用实现图解

    Jmeter跨线程组传值调用实现图解

    这篇文章主要介绍了Jmeter跨线程组传值调用实现图解,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友可以参考下
    2020-07-07

最新评论