【发布时间】:2021-08-05 15:43:29
【问题描述】:
我有发送消息的服务
@Service
class ExportTaskService {
@Autowired
private KafkaTemplate<String, Object> template;
public void exportNewTask(ImportTaskRequest req) {
template.send('my-topic-name', req)
}
}
我配置了bean:consumerFactory、producerFactory、kafkaTemplate(src/main/java)
如果我运行应用程序并执行metod -- 一切正常并在真正的消息代理中发送消息。
然后我需要 spring 测试,它使用 ExportTaskService.exportNewTask(request) 并等待来自同一主题的消息。
我的代码,但不起作用(我无法接收消息):
@RunWith(SpringRunner.class)
@SpringBootTest
@DirtiesContext
@TestPropertySource(locations="classpath:test.properties")
@EnableKafka
@EmbeddedKafka(
topics = "new-bitrix-leads", ports = 9092
)
public class ExportingLeadTests {
@Autowired
private EmbeddedKafkaBroker embeddedKafkaBroker;
@Autowired
ExportTaskService exportTaskService;
@Autowired
ConsumerFactory<String, Object> consumerFactory;
@Test
public void test() throws InterruptedException {
assert(embeddedKafkaBroker != null);
assert(exportTaskService != null);
Consumer<String, Object> consumer = consumerFactory.createConsumer();
consumer.subscribe(Collections.singletonList("new-bitrix-leads"));
exportTaskService.exportNewTask(ImportTaskRequest.builder()
.description("descr")
.title("title")
.build());
ConsumerRecords<String, Object> records = consumer.poll(Duration.ofSeconds(3));
assert (records.count() == 1);
}
}
我如何阅读此消息?我需要做什么 ?我没有想法...
大 TNX :) !!
【问题讨论】:
-
你看过spring-kafka github中已有的测试了吗?
标签: spring spring-boot apache-kafka spring-kafka