【问题标题】:Flink: Sharing state in CoFlatMapFunctionFlink:CoFlatMapFunction 中的共享状态
【发布时间】: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


【解决方案1】:

这是邮件列表中的答案...

CoFlatMapFunction 是否打算并行执行?

如果是,您需要某种方法来确定性地分配哪条记录 转到哪个并行实例。在某种程度上,CoFlatMapFunction 在模型和结果之间进行并行(分区)连接 会话窗口,因此您需要某种形式的键来选择哪个 分区元素去。这有意义吗?

如果不是,请尝试将其显式设置为并行度 1。

你好,斯蒂芬


所有人都可以只读访问的全局状态可以通过 广播()。

可供所有人读取和更新的全局状态是 当前不可用。一致的操作将是相当的 代价高昂,需要某种形式的分布式通信/共识。

相反,我鼓励您使用以下方法:

1) 如果您可以对状态进行分区,请使用 keyBy().mapWithState() - 那 本地化状态操作并使其非常快速。

2) 如果你的状态不是按键组织的,你的状态可能很 小,您也许可以使用非并行操作。

3) 如果某个操作更新了状态并且另一个操作访问了它, 您通常可以通过迭代和 CoFlatMapFunction 来实现它 (一侧是原始输入,另一侧是反馈输入)。

最终所有方法都本地化状态访问和修改, 如果可能的话,这是一个很好的模式。

你好,斯蒂芬

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2023-03-22
    • 1970-01-01
    • 2022-09-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-07
    相关资源
    最近更新 更多