【发布时间】:2023-03-07 06:06:01
【问题描述】:
我正在尝试为我实现 spring 云流 Kafka 流绑定器的拓扑实现测试代码。我正在使用功能样式,所以我想在商店中测试拓扑结果。但是当我调用时从商店返回 null 。你有什么想法吗?
这是我的测试代码;
@Log4j2
@RunWith(SpringRunner.class)
@SpringBootTest(
webEnvironment = SpringBootTest.WebEnvironment.NONE,
properties = {"server.port=0"})
public class ApplicationTests {
@ClassRule
public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, 1,
"input-topic");
private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka();
@Autowired
private QueryService queryService;
@Test
public void SimpleProcessorApplicationTest() {
//I'm producing data here
var result = queryService.getFromStore();
assert (actualResultSet.equals(result));
}
}
测试项目的配置文件;
spring:
cloud:
stream:
function:
definition: event1;event2;event3
bindings:
event1-in-0:
destination: input-topic
consumer:
timestampExtractorBeanName: eventTimeExtractor
dlqName: detail-dlq
event2-in-0:
destination: input-topic
consumer:
timestampExtractorBeanName: eventTimeExtractor
dlqName: detail2-dlq
event3-in-0:
destination: input-topic
consumer:
timestampExtractorBeanName: eventTimeExtractor
dlqName: detail3-dlq
kafka:
streams:
binder:
state-store-retry:
max-attempts: 2
backoff-period: 1000
replication-factor: 1
brokers: ${spring.embedded.kafka.brokers}
configuration:
commit.interval.ms: 10000
state.dir: state-store-test
application.server: 127.0.0.1:8080
default:
key:
serde: org.apache.kafka.common.serialization.Serdes$StringSerde
value:
serde: org.apache.kafka.common.serialization.Serdes$StringSerde
functions:
event1:
applicationId: aa-event4
event2:
applicationId: aa-event5
event3:
applicationId: aa-event6
deserialization-exception-handler: sendtodlq
是 QueryService 类;
@Log4j2
@Service
public class QueryService {
@Autowired
InteractiveQueryService interactiveQueryService;
public List<KeyValue<Integer, Result>> maxResults(){
List<KeyValue<Integer, Result>> allResult =new ArrayList<>();
final List<HostInfo> hostInfoList =
interactiveQueryService.getAllHostsInfo("m1-store");
for(HostInfo info: hostInfoList){
if(info.equals(interactiveQueryService.getCurrentHostInfo())){
log.info("Retrieving all key/value pairs from Local...");
allResult.addAll(getAllValues());
}
else{
log.info("Retrieving all key/value pairs from Remote...");
allResult.addAll(getAllValuesFromRemote(info));
}
}
return allResult;
}
private List<KeyValue<Integer, Result>> getAllValues() {
List<KeyValue<Integer, Result>> results = new ArrayList<>();
ReadOnlyKeyValueStore<Integer, Result> resultStore = interactiveQueryService.getQueryableStore(
"m1-store", QueryableStoreTypes.keyValueStore());
resultStore.all().forEachRemaining(results::add);
return results;
}
private List<KeyValue<Integer, Result>> getAllValuesFromRemote(HostInfo hostInfo){
String targetHost = String.format("http://%s:%d/dept/local", hostInfo.host(), hostInfo.port());
RestTemplate restTemplate = new RestTemplate();
List<KeyValue<Integer, Result>> result = restTemplate.getForObject(targetHost,List.class);
return result;
}
}
【问题讨论】:
-
它应该真正填充状态存储。你如何填充它?在测试方面,使用
EmbeddedKafka与使用真实集群没有任何不同。您是否尝试过调试您的测试并查看它是否确实被填充? -
是的,我看到它在物理上填充它。偏移量增加。查看测试函数中的注释。我也注释掉了,我调试的时候也看到了。
-
有没有机会在 GitHub 上分享这个项目,以便我们运行测试?
-
感谢您的大力支持!我已经在 GitHub 临时存储库上分享了。这是回购; github.com/kadiralan/kafka-test
标签: spring-boot apache-kafka spring-cloud apache-kafka-streams spring-cloud-stream