【问题标题】:Spring Kafka. Not started EmbeddedKafkaBroker春天的卡夫卡。未启动 EmbeddedKafkaBroker
【发布时间】:2020-03-25 23:47:50
【问题描述】:

我正在编写 Kafka Broker 和 Consumer 代码以捕获来自应用程序的消息。尝试从 Consumer 获取消息时,出现错误

java.net.ConnectException: Connection refused: no further information
    at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method)
    at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:717)
    at org.apache.kafka.common.network.PlaintextTransportLayer.finishConnect(PlaintextTransportLayer.java:50)
    at org.apache.kafka.common.network.KafkaChannel.finishConnect(KafkaChannel.java:216)
    at org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:531)
    at org.apache.kafka.common.network.Selector.poll(Selector.java:483)
    at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:540)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:262)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:233)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:212)
    at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureCoordinatorReady(AbstractCoordinator.java:230)
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:444)
    at org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1267)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1231)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1211)
    at org.springframework.kafka.test.utils.KafkaTestUtils.getRecords(KafkaTestUtils.java:303)
    at org.springframework.kafka.test.utils.KafkaTestUtils.getRecords(KafkaTestUtils.java:280)

在应用端(Producer),也存在连接错误

2020-03-25 12:29:33.689  WARN 25786 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient   : [Producer clientId=producer-1, transactionalId=tx0] Connection to node -1 (<here broker hostname>:9092) could not be established. Broker may not be available.

我的项目有以下依赖:

compile "org.springframework.kafka:spring-kafka-test:2.4.4.RELEASE"
compile "org.springframework.kafka:spring-kafka:2.4.4.RELEASE"

我的 Kafka 经纪人的代码

public class KafkaServer {

    private static final String BROKERPORT = "9092";
    private static final String BROKERHOST = "localhost";
    public static final String TOPIC1 = "fss-fsstransdata";
    public static final String TOPIC2 = "fss-fsstransscores";
    public static final String TOPIC3 = "fss-fsstranstimings";
    public static final String TOPIC4 = "fss-fssdevicedata";
    @Getter
    private Consumer<String, String> consumer;

    private EmbeddedKafkaBroker embeddedKafkaBroker;

    public void run() {

        String[] topics = {TOPIC1, TOPIC2, TOPIC3, TOPIC4};

        this.embeddedKafkaBroker = new EmbeddedKafkaBroker(
                1,
                false,
                1,
                topics
        ).kafkaPorts(BROKERPORT);

        Map<String, Object> configs = new HashMap<>(KafkaTestUtils.consumerProps("consumer", "false", this.embeddedKafkaBroker));
        this.consumer = new DefaultKafkaConsumerFactory<>(configs, new StringDeserializer(), new StringDeserializer()).createConsumer();

        this.consumer.subscribe(Arrays.asList(topics));
    } 
}

请帮助处理这种情况。我不擅长 kafka 架构以及如何在 Spring 上实现它。

【问题讨论】:

  • 顺便说一句,Kafka 消费者和 Kafka 经纪人不一样。

标签: java apache-kafka spring-kafka


【解决方案1】:

如果您使用的是 spring,那么您需要使用 @EmbeddedKafka 注释您的 bean,然后在 EmbeddedKafkaBroker 上使用 @Autowire

嵌入式kafka注解配置示例:

@EmbeddedKafka(
    partitions = 1, 
    controlledShutdown = false,
    brokerProperties = {// place your proerties here
})

我要做的是创建一个 spring bean KafkaServerConfig 并将我所有的配置和 bean 构造逻辑放在里面。

PS:需要注意的是 EmbeddedKafkaBroker 用于单元测试。

【讨论】:

  • 你好,亚历山大,谢谢你的回答!是否可以在不使用注释和 DI 的情况下完成所有这些操作?我不想创建一个在代码中手动创建所有需要的对象的弹簧上下文。 p.s.我需要一个 Kafka 存根来验证从向 kafka 发送消息的应用程序发送的消息的正确性,仅此而已
【解决方案2】:

EmbeddedKafkaBroker 旨在从 Spring 应用程序上下文或由 JUnit4 @Rule@ClassRule 或 JUnit5 Condition 使用。

要在这些环境之外使用它,您必须调用 afterPropertiesSet() 来初始化它并调用 destroy() 来关闭它。

【讨论】:

  • 你好,加里!谢谢你的回答。我在 EmbeddedKafkaBroker 创建后使用 afterPropertiesSet(),问题“java.net.ConnectException:连接被拒绝”消失了。我的存根仍然没有收到来自应用程序的消息,也没有错误,但这是另一回事)
  • 欢迎堆栈溢出!见stackoverflow.com/help/someone-answers
  • 如果你在测试用例中使用它,最常见的没有收到消息的原因是因为它是在消费者被分配分区之前发送的。将ConsumerConfig.AUTO_OFFSET_RESET_CONFIG 设置为earliest - 默认情况下为latest,因此当您被分配时,您将不会获得已在主题中的任何记录。
  • 不,我的问题是生产者无法连接到kafka broker,我配置错误。我决定尝试在其中一个 linux 测试台上部署 kafka 服务器并向其发送消息)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-12-17
  • 2017-02-08
  • 2018-02-10
  • 1970-01-01
  • 2015-06-18
相关资源
最近更新 更多