【问题标题】:How do I resolve an error in my configuration when trying to unit test a Spring Kafka Consumer?尝试对 Spring Kafka Consumer 进行单元测试时,如何解决配置中的错误?
【发布时间】:2018-07-21 06:11:28
【问题描述】:

代码位置

我认为可能有很多模块可以使这个问题看起来干净,所以这里是 repo。我希望包括所有必要的组件。 https://github.com/ewingian/RestCalculator

问题

我正在学习编写 Kafka 服务,这个过程包括为生产者和消费者学习单元测试。遵循有关与消费者一起设置单元测试的教程。当我运行测试时,我收到一个类配置错误。

错误

9:34:59.547 [main] WARN kafka.server.BrokerMetadataCheckpoint - No meta.properties file under dir /tmp/kafka-8346130278143417083/meta.properties
09:34:59.567 [main] ERROR kafka.server.KafkaServer - [Kafka Server 0], Fatal error during KafkaServer startup. Prepare to shutdown
java.lang.NoClassDefFoundError: org/apache/kafka/common/network/LoginType
    at kafka.network.Processor.<init>(SocketServer.scala:406)
    at kafka.network.SocketServer.newProcessor(SocketServer.scala:141)
    at kafka.network.SocketServer$$anonfun$startup$1$$anonfun$apply$1.apply$mcVI$sp(SocketServer.scala:94)
    at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:160)
    at kafka.network.SocketServer$$anonfun$startup$1.apply(SocketServer.scala:93)
    at kafka.network.SocketServer$$anonfun$startup$1.apply(SocketServer.scala:89)
    at scala.collection.Iterator$class.foreach(Iterator.scala:893)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1336)
    at scala.collection.MapLike$DefaultValuesIterable.foreach(MapLike.scala:206)
    at kafka.network.SocketServer.startup(SocketServer.scala:89)
    at kafka.server.KafkaServer.startup(KafkaServer.scala:219)
    at kafka.utils.TestUtils$.createServer(TestUtils.scala:120)
    at kafka.utils.TestUtils.createServer(TestUtils.scala)
    at org.springframework.kafka.test.rule.KafkaEmbedded.before(KafkaEmbedded.java:154)
    at org.junit.rules.ExternalResource$1.evaluate(ExternalResource.java:46)
    at org.junit.rules.RunRules.evaluate(RunRules.java:20)
    at org.junit.runners.ParentRunner.run(ParentRunner.java:363)
    at org.springframework.test.context.junit4.SpringJUnit4ClassRunner.run(SpringJUnit4ClassRunner.java:191)
    at org.junit.runner.JUnitCore.run(JUnitCore.java:137)
    at com.intellij.junit4.JUnit4IdeaTestRunner.startRunnerWithArgs(JUnit4IdeaTestRunner.java:68)
    at com.intellij.rt.execution.junit.IdeaTestRunner$Repeater.startRunnerWithArgs(IdeaTestRunner.java:51)
    at com.intellij.rt.execution.junit.JUnitStarter.prepareStreamsAndStart(JUnitStarter.java:242)
    at com.intellij.rt.execution.junit.JUnitStarter.main(JUnitStarter.java:70)
Caused by: java.lang.ClassNotFoundException: org.apache.kafka.common.network.LoginType
    at java.net.URLClassLoader.findClass(URLClassLoader.java:381)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
    at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:335)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
    ... 23 common frames omitted

图为我的目录结构

当我查看 IDEA 项目结构设置中的库列表时,我看到了org.apache.kafka:kafka-clients:0.11.0.0;但是我无法导入缺少的模块,我知道它是 kafka-clients 的一部分。 (org/apache/kafka/common/network/LoginType)

问题

以前有人遇到过这个错误吗?我是否错误配置了我的 gradle 文件?我的项目目录是否正确设置以有效地我可能缺少什么?执行 Kafka 单元测试?我还没有找到关于 LoginType 的太多信息,但会继续搜索。

这里是 gradle 构建文件的副本:

buildscript {
    repositories {
        mavenCentral()
    }misconfigured
    dependencies {
        classpath("org.springframework.boot:spring-boot-gradle-plugin:1.5.10.RELEASE")
    }
}

apply plugin: 'java'
apply plugin: 'eclipse'
apply plugin: 'idea'
apply plugin: 'org.springframework.boot'

jar {
    baseName = 'calculator'
    version =  '0.1.0'
}

repositories {
    mavenCentral()
}

sourceCompatibility = 1.8
targetCompatibility = 1.8

dependencies {
    compile("org.springframework.boot:spring-boot-starter-web")
    testCompile("org.springframework.boot:spring-boot-starter-test")
    compile("org.springframework.kafka:spring-kafka:1.3.2.RELEASE")
    testCompile("org.springframework.kafka:spring-kafka-test")
}

task wrapper(type: Wrapper) {
    gradleVersion = '2.3'
}

如果我可能需要在此问题中包含其他任何内容,请告诉我。谢谢

单元测试代码

package com.calculator;

/**
 * Created by ian on 2/9/18.
 */
import com.calculator.kafka.services.*;

import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertThat;
import static org.springframework.kafka.test.assertj.KafkaConditions.key;
import static org.springframework.kafka.test.hamcrest.KafkaMatchers.hasValue;

import java.util.Map;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.junit.After;
import org.junit.Before;
import org.junit.ClassRule;
import org.junit.Test;

import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.listener.config.ContainerProperties;
import org.springframework.kafka.test.rule.KafkaEmbedded;
import org.springframework.kafka.test.utils.KafkaTestUtils;

import org.springframework.test.context.junit4.SpringRunner;

@RunWith(SpringRunner.class)
@SpringBootTest
public class KafkaTest {

//    private static final Logger LOGGER = LoggerFactory.getLogger(KafkaTest.class);
    private static final String TEMPLATE_TOPIC = "input";
    private static String SENDER_TOPIC = "input";
    @ClassRule
    public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEMPLATE_TOPIC);

    private KafkaMessageListenerContainer<String, Integer> container;

    private BlockingQueue<ConsumerRecord<String, Integer>> records;

    @Autowired
    private KafkaConsumer consumer;

    @Before
    public void setUp() {
        // Set up the consumer properties
        Map<String, Object> integerProperties = KafkaTestUtils.consumerProps("jsa-group", "false", embeddedKafka);

        // create a Kafka consumer factory
        DefaultKafkaConsumerFactory<String, Integer> consumerFactory = new DefaultKafkaConsumerFactory<String, Integer>(integerProperties);

        // set the topic that needs to be consumed
        ContainerProperties containerProperties = new ContainerProperties(SENDER_TOPIC);

        // create a Kafka MessageListenerContainer
        container = new KafkaMessageListenerContainer<>(consumerFactory, containerProperties);

        // setup a Kafka message listener
        container.setupMessageListener(new MessageListener<String, Integer>() {
            @Override
            public void onMessage(ConsumerRecord<String, Integer> record) {
//                LOGGER.debug("test-listener received message='{}'", record.toString());
                records.add(record);
            }
        });

        // start the container and underlying message listener
        container.start();
    }

    @After
    public void tearDown() {
        // stop the container
        container.stop();
    }

    @Test
    public void testTemplate() throws Exception {
        // send the message
        String greeting = "Hello Spring Kafka Sender!";
        Integer i1 = 12;
        consumer.processMessage(i1);

        // check that the message was received
        ConsumerRecord<String, Integer> received = records.poll(10, TimeUnit.SECONDS);
        // Hamcrest Matchers to check the value
        assertThat(received, hasValue(i1));
        // AssertJ Condition to check the key
        assertThat(received).has(key(null));
    }
}

【问题讨论】:

  • 嘿@pvpkiran 我想我错过了连接。除非设置 CountDownLatch,并将偏移量设置为最早以某种方式否定 org/apache/kafka/common/network/LoginType

标签: spring unit-testing intellij-idea apache-kafka


【解决方案1】:

以下是我发现的摘要:

  1. 我没有使用最新的 Spring Boot 版本。按照教程,我使用了 1.5.10.RELEASE。没有错,但它与 spring kafka 存在兼容性问题。
  2. 我尝试使用最新版本的 spring kafka 和 spring boot 1.5.10.RELEASE。不断收到处理未找到的某些类的错误。不得不将 Kafka 的版本降低到 1.3.2.RELEASE
  3. 此配置允许我运行 Spring Boot 应用程序,但单元测试失败,因此出现此堆栈溢出问题。
  4. 我试图重写我的 gradle 文件以使用更新的 spring boot 和 kafka,但失败了,我认为存储库不好。

解决方案

终于去了spring的网站并使用了他们的项目生成器。引入了最新的春季 Kafka 和 Boost。这给了我一个新的 build.gradle 和正确的 repos。手动添加spring-kafka-test,单元测试成功。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-30
    • 2018-11-10
    • 1970-01-01
    • 2019-10-25
    • 2020-03-26
    相关资源
    最近更新 更多