【问题标题】:Failed to create or validate data directory for EmbeddedKafkaBroker无法为 EmbeddedKafkaBroker 创建或验证数据目录
【发布时间】:2020-02-10 12:10:38
【问题描述】:

我正在尝试使用@EmbeddedKafka Annotation 运行一个简单的单元测试。 作为参考,我正在关注以下春季文档 https://docs.spring.io/spring-kafka/reference/html/#embedded-kafka-annotation

@RunWith(SpringRunner.class)
@DirtiesContext
@EmbeddedKafka(brokerProperties = "log.dir=/kafka-logs", partitions = 1,
    topics = {
        "dare_policy_created"})
@Slf4j
public class ConsumerTest {

@Autowired
  private EmbeddedKafkaBroker embeddedKafka;

@Test
  public void someTest() {
    Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("testGroup", "true", this.embeddedKafka);
    consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    ConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<>(consumerProps);
    Consumer<Integer, String> consumer = cf.createConsumer();
    this.embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "dare_policy_created");
    ConsumerRecords<Integer, String> replies = KafkaTestUtils.getRecords(consumer);
    //assertThat(replies.count()).isGreaterThanOrEqualTo(1);
  }
}

我试图定义 log.dir @EmbeddedKafka(brokerProperties = "log.dir= ") 因为我在运行测试时遇到错误。

我试过了:

  • log.dir=/kafka-logs
  • log.dir=real_path_to_my_project/kafka-logs
  • ...

但是每次我运行测试时都会出现这个错误:

kafka.server.LogDirFailureChannel.error - Failed to create or validate data directory /kafka-logs java.io.IOException: Failed to load /kafka-logs during broker startup

kafka.log.LogManager.fatal - Shutdown broker because none of the specified log dirs from /kafka-logs can be created or validated

【问题讨论】:

  • 我也有这个问题。通过将 meta.properties 文件添加到日志目录,我更进一步。 Kafka 尝试启动但随后失败:java.lang.NoSuchMethodError: org.apache.kafka.common.network.ChannelBuilders.serverChannelBuilder(Lorg/apache/kafka/common/network/ListenerName;ZLorg/apache/kafka/common/security /auth/SecurityProtocol;Lorg/apache/kafka/common/config/AbstractConfig;Lorg/apache/kafka/common/security/authenticator/CredentialCache;Lorg/apache/kafka/common/security/token/delegation/internals/DelegationTokenCache;) Lorg/apache/kafka/common/network/ChannelBuilder;
  • 我也遇到了同样的错误(+1),你找到解决办法了吗?

标签: spring-boot junit apache-kafka embedded-kafka


【解决方案1】:

我能够通过删除对 kafka-client 的显式依赖来解决该问题。 我的 pom 中有以下依赖项

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.3.0</version>
</dependency>

【讨论】:

    猜你喜欢
    • 2017-01-24
    • 2018-12-13
    • 1970-01-01
    • 2016-06-04
    • 1970-01-01
    • 1970-01-01
    • 2020-01-20
    • 2015-05-26
    • 1970-01-01
    相关资源
    最近更新 更多