【发布时间】:2019-03-03 01:28:16
【问题描述】:
我在使用来自 https://mvnrepository.com/artifact/org.springframework.kafka/spring-kafka-test/2.1.10.RELEASE 的 KafkaEmbedded 时遇到问题
我正在使用KafkaEmbedded 创建用于测试生产者/消费者管道的 Kafka 代理。这些生产者/消费者是来自 kafka-clients 的标准客户。我没有使用 Spring Kafka 客户端。
一切正常,代码运行良好,但我必须使用 KafkaEmbedded 中的 consumeFromEmbeddedTopics() 方法才能使消费者正常工作。如果我不使用此方法,消费者将不会收到任何消息。
这个方法有两个问题:首先,它需要KafkaConsumer 作为参数(我不想在类中公开它)并且调用这个方法会在对象调用 poll 时给出ConcurrentModificationException @Scheduled.
我使用的是auto.offset.reset 属性,所以这是另一回事。
我的问题是:如何在不调用这些 consumeFromEmbeddedTopics() 方法的情况下正确使用来自 KafkaEmbedded 的记录?
【问题讨论】:
标签: apache-kafka spring-test spring-kafka