【问题标题】:How to test state store in apache kafka using spring cloud stream?如何使用 Spring Cloud Stream 测试 apache kafka 中的状态存储?
【发布时间】: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


【解决方案1】:

我发现了几个问题。您在配置中的application.server 属性指向localhost:8080。将其更改为 ${spring.embedded.kafka.brokers}。此外,在您的测试中,您发送到错误的主题名称:template.setDefaultTopic("stage.dashboard.iddaa.coupondetail");。 您在配置中将目的地定义为input-topic。在进行这些更改之后,我能够更进一步,但是,您的 aggregate 呼叫根本没有被调用。我没有机会对此进行调试,但我怀疑这与您的窗口逻辑有关。基本上,您的处理器没有将任何内容放入商店agg-store,这就是对商店的查询返回空的原因。看看为什么会这样。希望这些是您进一步分类的一些指示。

【讨论】:

  • 非常感谢!我非常感谢您的支持。我用“输入主题”名称替换了主题名称。但是,据我们所见,我忘了用输入主题替换它。我会听取您的建议并进行更新。再次感谢!
  • 不,我无法修复它。我看到了一个内部异常。我认为应用程序服务器在我调用它时没有运行。你有什么主意吗?这是例外;无法获取状态存储 agg-store,因为流线程是 PARTITIONS_ASSIGNED,而不是 RUNNING
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多