【问题标题】:how to send json data stream to multiple topics in kafka based on input fields如何根据输入字段将json数据流发送到kafka中的多个主题
【发布时间】:2020-03-02 12:38:25
【问题描述】:

我必须使用来自 kafka 流的 json 数据并将其发送到不同的主题(应用程序 ID 和实体的不同组合)以供进一步使用。
主题名称:

    app1.entity1
    app1.entity2
    app2.entity1
    app2.entity2

Json 数据

    [
        {
            "appId": "app1",
            "entity": "entity1",
            "extractType": "txn",
            "status": "success",
            "fileId": "21151235"
        },
        {
            "appId": "app1",
            "entity": "entity2",
            "extractType": "txn",
            "status": "fail",
            "fileId": "2134234123"
        },
        {
            "appId": "app2",
            "entity": "entity3",
            "extractType": "payment",
            "status": "success",
            "fileId": "2312de23e"
        },
        {
            "appId": "app2",
            "entity": "entity3",
            "extractType": "txn",
            "status": "fail",
            "fileId": "asxs3434"
        }
    ]

TestInput.java

        private String appId;           
        private String entity ;             
        private String extractType;         
        private String status;          
        private String fileId; 

        setter/gtter

SpringBootConfig.java

      @Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME)
      public KafkaStreamsConfiguration kStreamsConfigs(KafkaProperties kafkaProperties) {
          Map<String, Object> config = new HashMap<>();
          config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
          config.put(StreamsConfig.APPLICATION_ID_CONFIG, kafkaProperties.getClientId());
          config.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
          config.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, new JsonSerde<>(TestInput.class).getClass());
          config.put(JsonDeserializer.DEFAULT_KEY_TYPE, String.class);
          config.put(JsonDeserializer.DEFAULT_VALUE_TYPE, TestInput.class);
          return new KafkaStreamsConfiguration(config);
      }

      @Bean
      public KStream<String, TestInput> kStream(StreamsBuilder kStreamBuilder) {
          KStream<String, TestInput> stream = kStreamBuilder.stream(inputTopic);
                 // how to form key , group records and send to different topics
          return stream;
      }

我搜索了很多,但没有找到任何将数据动态发布到主题的内容。请高手帮忙

【问题讨论】:

    标签: spring-boot apache-kafka-streams spring-kafka


    【解决方案1】:

    使用stream.branch()

    https://www.confluent.io/blog/putting-events-in-their-place-with-dynamic-routing/

    接下来,让我们修改需求。每个微服务不应处理流中的所有事件,而应仅对相关事件的子集采取行动。处理此要求的一种方法是让微服务订阅所有事件的原始流,检查每条记录,然后只对它关心的事件采取行动,而丢弃其余的。但是,根据应用程序的不同,这可能是不可取的或资源密集型的。

    一种更简洁的方法是为服务提供一个单独的流,该流仅包含微服务关心的相关事件子集。为此,流应用程序可以使用方法 KStream#branch() 将原始事件流分支到不同的子流中。这会产生新的 Kafka 主题,因此微服务可以直接订阅其中一个分支流。

    ...

    【讨论】:

    • 加里,有没有任何示例或链接可供我参考以获取有关实施的一些提示
    • 哪一个?自定义处理器还是使用 spring-kafka?
    • 自定义处理器
    • 谢谢加里。你的指针真的帮助我找到了解决方案。 stackoverflow.com/questions/48950580/…
    猜你喜欢
    • 2019-07-12
    • 1970-01-01
    • 2021-06-03
    • 2018-08-03
    • 2017-05-21
    • 1970-01-01
    • 2021-03-01
    • 2020-07-17
    • 2023-03-29
    相关资源
    最近更新 更多