【问题标题】:Apache flink broadcast state gets flushedApache flink 广播状态被刷新
【发布时间】:2019-05-22 21:02:29
【问题描述】:

我正在使用广播模式连接两个流并从一个流读取数据。代码是这样的

case class Broadcast extends BroadCastProcessFunction[MyObject,(String,Double), MyObject]{
  override def processBroadcastElement(in2: (String, Double), 
                                       context: BroadcastProcessFunction[MyObject, (String, Double), MyObject]#Context,
                                       collector:Collector[MyObject]):Unit={
    context.getBroadcastState(broadcastStateDescriptor).put(in2._1,in2._2)
  }

  override def processElement(obj: MyObject,
                            readOnlyContext:BroadCastProcessFunction[MyObject, (String,Double), 
                            MyObject]#ReadOnlyContext, collector: Collector[MyObject]):Unit={
    val theValue = readOnlyContext.getBroadccastState(broadcastStateDesriptor).get(obj.prop)
    //If I print the context of the state here sometimes it is empty.
    out.collect(MyObject(new, properties, go, here))
  }
}

状态描述符:

val broadcastStateDescriptor: MapStateDescriptor[String, Double) = new MapStateDescriptor[String, Double]("name_for_this", classOf[String], classOf[Double])

我的执行代码如下所示。

val streamA :DataStream[MyObject] = ... 
val streamB :DataStream[(String,Double)] = ... 
val broadcastedStream = streamB.broadcast(broadcastStateDescriptor)

streamA.connect(streamB).process(new Broadcast)

问题出在processElement 函数中,状态有时为空,有时不是。状态应该始终包含数据,因为我不断地从一个我知道它有数据的文件中流式传输。我不明白为什么它正在刷新状态并且我无法获取数据。

我尝试在将数据放入状态之前和之后在processBroadcastElement中添加一些打印,结果如下

0 - 1
1 - 2 
2 - 3 
.. all the way to 48 where it resets back to 0

更新: 我注意到的是,当我减少流执行上下文的超时值时,结果会好一些。当我增加它时,地图总是空的。

env.setBufferTimeout(1) //better results 
env.setBufferTimeout(200) //worse result (default is 100)

【问题讨论】:

    标签: scala streaming state apache-flink


    【解决方案1】:

    每当在 Flink 中连接两个流时,您无法控制 Flink 将事件从两个流传递到您的用户函数的时间。因此,例如,如果存在可从流 A 处理的事件,以及可从流 B 处理的事件,则接下来可能会处理其中一个事件。您不能期望 broadcastedStream 以某种方式优先于其他流。

    根据您的要求,您可以采用各种策略来应对这两个流之间的竞争。例如,您可以使用 KeyedBroadcastProcessFunction 并使用其 applyToKeyedState 方法在新的广播事件到达时迭代所有现有的键控状态。

    【讨论】:

    • 但是状态怎么变空了呢?如果我让流运行几分钟,那么我希望状态一直保持满,对吗?但事实并非如此。看起来状态被重置了。
    • 如果您分享更多代码,可能会了解您的应用在做什么。
    • 广播状态是一个Map——当你说它为空时,你的意思是没有你希望有值的特定键的值,或者没有有值的键什么?
    • 我的意思是没有带值的键。地图完全是空的。这就是让我好奇的原因。地图上有有效期吗?地图在读取广播流后不应该一直是满的吗?
    • 我注意到当我减少流执行上下文的数字超时值时,结果会好一些。当我增加它时,地图总是空的。 env.setBufferTimeout(1) //更好的结果 env.setBufferTimeout(200) //更差的结果
    【解决方案2】:

    正如大卫所说,这项工作可能会重新开始。我禁用了检查点,因此我可以看到任何可能引发的异常,而不是 flink 静默失败并重新启动作业。

    原来是在尝试解析文件时出错。因此作业不断重启,因此状态为空,flink 不断消耗流。

    【讨论】:

    • 很高兴听到您发现了问题。顺便说一下,查看日志(或检查点和正常运行时间/重启指标)应该会发现这一点。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多