【发布时间】:2015-11-17 11:23:27
【问题描述】:
有点卡住CoFlatMapFunction。如果我将它放在窗口之前的DataStream 上似乎可以正常工作,但如果放在窗口的“应用”功能之后则会失败。
我正在测试两个流,flatMap1 上的主要“功能”不断摄取数据,flatMap2 上的控制流“模型”根据要求更改模型。
我可以在flatMap2 中正确设置并查看 b0/b1,但 flatMap1 总是看到 b0 和 b1,因为在初始化时设置为 0。
我在这里遗漏了什么明显的东西吗?
public static class applyModel implements CoFlatMapFunction<Features, Model, EnrichedFeatures> {
private static final long serialVersionUID = 1L;
Double b0;
Double b1;
public applyModel(){
b0=0.0;
b1=0.0;
}
@Override
public void flatMap1(Features value, Collector<EnrichedFeatures> out) {
System.out.print("Main: " + this + "\n");
}
@Override
public void flatMap2(Model value, Collector<EnrichedFeatures> out) {
System.out.print("Old Model: " + this + "\n");
b0 = value.getB0();
b1 = value.getB1();
System.out.print("New Model: " + this + "\n");
}
@Override
public String toString(){
return "CoFlatMapFunction: {b0: " + b0 + ", b1: " + b1 + "}";
}
}
【问题讨论】:
-
你在窗口
apply函数中做了什么?也许您可以与我们分享各自的代码。 -
for(原始值:值){ if(value.getTs() > end_ts) end_ts = value.getTs(); if (value.getTs()
-
我会检查是否可以重现您的问题。
-
并行实例是我的问题。必须确保两个流的键控方式相同,将我的 CoFlatMapFunction 的并行度设置为 1。感谢 Stephan Ewen MORE INFORMATION
-
很高兴听到。您想在此处发布 Stephan 的答案,以便其他有相同问题的人可以轻松找到解决方案吗?
标签: apache-flink flink-streaming