【问题标题】:Kafka Stream producing custom list of messages based on certain conditionsKafka Stream 根据特定条件生成自定义消息列表
【发布时间】:2019-11-15 18:47:01
【问题描述】:

我们有以下流处理要求。

Source Stream -> 
 transform(condition check - If (true) then generate MULTIPLE ADDITIONAL messages else just transform the incoming message) ->
 output kafka topic

Example:
If condition is true for message B(D,E,F are the additional messages produced)
A,B,C -> A,D,E,F,C -> Sink Kafka Topic
If condition is false     
A,B,C -> A,B,C -> Sink Kafka Topic

有没有办法在 Kafka 流中实现这一点?

【问题讨论】:

    标签: apache-kafka kafka-consumer-api apache-kafka-streams apache-kafka-connect confluent-platform


    【解决方案1】:

    您可以使用flatMap()flatMapValues() 方法。这些方法获取一条记录并产生零、一条或多条记录。

    flatMap() 可以修改键、值及其数据类型,而flatMapValues() 保留原始键并更改值和值数据类型。

    这是一个伪代码示例,考虑到新消息“C”、“D”、“E”将有一个新密钥。

    KStream<byte[], String> inputStream = builder.stream("inputTopic");
    KStream<byte[], String> outStream = inputStream.flatMap( 
               (key,value)->{
                List<KeyValue<byte[], String>> result = new LinkedList<>();  
                    // If message value is "B". Otherwise place your condition based on data     
                    if(value.equalsTo("B")){ 
                          result.add(KeyValue.pair("<new key for message C>","C"));
                          result.add(KeyValue.pair("<new key for message D>","D"));
                          result.add(KeyValue.pair("<new key for message E>","E"));
    
                     }else{
                             result.add(KeyValue.pair(key,value));
                     }
                return result;
    });
    outStream.to("sinkTopic");
    

    您可以阅读更多相关信息: https://docs.confluent.io/current/streams/developer-guide/dsl-api.html#streams-developer-guide-dsl-transformations-stateless

    【讨论】:

    • 如果我们想有条件地生产怎么办?假设我做了一些改造和改造结果的基础,我想决定是否推送到另一个主题?
    猜你喜欢
    • 1970-01-01
    • 2018-07-18
    • 2017-01-30
    • 1970-01-01
    • 2015-10-09
    • 2018-03-09
    • 1970-01-01
    • 1970-01-01
    • 2017-11-18
    相关资源
    最近更新 更多