【问题标题】:Precessing two DataStream<String> simultaneously ,to find one DataStream contains values from other DataStream in flink?同时处理两个 DataStream<String> ,找到一个 DataStream 包含来自 flink 中其他 DataStream 的值?
【发布时间】:2020-12-10 19:04:24
【问题描述】:

假设我有两个DataStream&lt;String&gt;,我收到了来自 Kafka 的流,经过一些处理,我得到了这两个流。

DataStream<String> A contains values {id1_id2 , id3_id4, id99_id0, id15_id3,id11_id5....}

DataStream<String> B contains values {id2, id3,id5...}

是否可以对DataStream A进行一些处理,以便将值输出到另一个中

DataStream<String> C ={id1, id3, id15, id11}

因此,B 中存在的所有值都将与 A 相交。我已尝试使用 processElement() 和 RichCoFlatMapFunction,但它不起作用。

public class MatchAggregator
        extends RichCoFlatMapFunction<String, String, Tuple1<String>> {

    private ValueState<String> doubleState;
    private ValueState<String> singleState;

    @Override
    public void open(Configuration config) {

        doubleState = getRuntimeContext().getState(new ValueStateDescriptor<>("doubleEvents",String.class));
        singleState = getRuntimeContext().getState(new ValueStateDescriptor<>("singleEvents",String.class));
    }
    
    @Override
    public void flatMap1(String s, Collector<Tuple1<String>> collector) throws Exception {
        String single = singleState.value();
       //this is outputting null.
        System.out.println(single);
      //s is also null
        if(single.contains(s)){
            String replaceNumber = single.replace(s,"");
            String replaceEmp = replaceNumber.replace("_","");
            single.clear();
            collector.collect(Tuple1.of(replaceEmp));
        }else {
            personContactState.update(s);
        }
    }

    @Override
    public void flatMap2(String s, Collector<Tuple1<String>> collector) throws Exception {

        
    }
}

我正在使用两个数据流,例如:

DataStream<Tuple1<String>> match = A.connect(B).flatMap(new MatchAggregator());

match.print();

【问题讨论】:

  • 来自 Flink 文档的 This tutorial 和来自 Flink 培训的 this exercise 将教您完成这项工作所需了解的知识。
  • 我已将 Rich 函数添加为您提供的教程,但出现 Nullpointer 异常。

标签: java apache-kafka apache-flink flink-streaming


【解决方案1】:

RichCoFlatMapFunction 的确切行为将取决于您如何键入两个连接的流。

String single = singleState.value() 将检索以前存储的任何值为相同的键作为传入String s 的键。在您共享的代码中,update 永远不会在 singleState 上调用,因此 singleState.value() 将始终为空。

【讨论】:

  • 感谢您的教程,它正在工作。我会发布答案,所以将来任何人都可以使用它。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-12-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多