java - 这种在Spring Boot应用程序中启动无限循环的方式有什么问题吗?

标签 java spring spring-boot apache-kafka

我有一个 Spring Boot 应用程序,它需要处理一些 Kafka 流数据。我向将在启动时运行的 CommandLineRunner 类添加了一个无限循环。里面有一个可以被唤醒的Kafka消费者。我用 Runtime.getRuntime().addShutdownHook(new Thread(consumer::wakeup)); 添加了一个关闭钩子(Hook)。我会遇到什么问题吗?在 Spring 中是否有更惯用的方法来执行此操作?我应该改用 @Scheduled 吗?下面的代码去除了特定的 Kafka 实现内容,但在其他方面是完整的。

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;

import java.time.Duration;
import java.util.Properties;


    @Component
    public class InfiniteLoopStarter implements CommandLineRunner {

        private final Logger logger = LoggerFactory.getLogger(this.getClass());

        @Override
        public void run(String... args) {
            Consumer<AccountKey, Account> consumer = new KafkaConsumer<>(new Properties());
            Runtime.getRuntime().addShutdownHook(new Thread(consumer::wakeup));

            try {
                while (true) {
                    ConsumerRecords<AccountKey, Account> records = consumer.poll(Duration.ofSeconds(10L));
                    //process records
                }
            } catch (WakeupException e) {
                logger.info("Consumer woken up for exiting.");
            } finally {
                consumer.close();
                logger.info("Closed consumer, exiting.");
            }
        }
    }

最佳答案

我不确定你是否会在那里遇到任何问题,但它有点脏 - Spring 对使用 Kafka 有很好的内置支持,所以我会倾向于它(网络上有很多关于它的文档,但一个不错的是:https://www.baeldung.com/spring-kafka)。

您需要以下依赖项:

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
    <version>2.2.2.RELEASE</version>
</dependency>

配置非常简单,只需将 @EnableKafka 注解添加到配置类,然后设置 Listener 和 ConsumerFactory bean

配置完成后,您可以按如下方式轻松设置消费者:

@KafkaListener(topics = "topicName")
public void listenWithHeaders(
  @Payload String message, 
  @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition) {
      System.out.println("Received Message: " + message"+ "from partition: " + partition);
}

关于java - 这种在Spring Boot应用程序中启动无限循环的方式有什么问题吗?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/54194780/

相关文章:

java - 指定调用的数组长度

java - 提取 DatabaseMetaData 时出错;嵌套异常是 com.mysql.jdbc.exceptions.jdbc4。连接关闭后不允许进行任何操作

java - 国家边界坐标数据

Java公私钥解密问题

java - WebLogic 集成的替代方案

java - JDBI org.skife.jdbi.v2.exceptions.UnableToCreateStatementException : Exception while binding

java - 在 Thymeleaf 中,当显示属于一个对象的所有属性时,其中一个属性是其他对象的列表,我该如何显示该列表?

java - JPA:删除孤孙元素时的问题

java - 拦截 Rest Controller 响应

java - 为什么 spring 框架不将我的 YAML 列表绑定(bind)到实体?