【问题标题】:How do I handle out-of-order events with Apache flink?如何使用 Apache flink 处理乱序事件?
【发布时间】:2021-01-31 14:53:06
【问题描述】:

为了测试流处理和 Flink,我给自己提出了一个看似简单的问题。我的数据流由粒子的xy 坐标以及记录位置的时间t 组成。我的目标是用特定粒子的速度注释这些数据。所以流可能看起来像这样。

<timestamp:Long> <particle_id:String> <x:Double> <y:Double>

1612103771212 p1 0.0 0.0
1612103771212 p2 0.0 0.0
1612103771213 p1 0.1 0.1
1612103771213 p2 -0.1 -0.1
1612103771214 p1 0.1 0.2
1612103771214 p2 -0.1 -0.2
1612103771215 p1 0.2 0.2
1612103771215 p2 -0.2 -0.2

现在无法保证事件会按顺序到达,即 1612103771213 p2 -0.1 -0.1 可能会在 10ms 之前到达 1612103771212 p2 0.0 0.0

为简单起见,可以假设任何迟到的数据都将在早期数据的100ms 内到达。

我承认我是流处理和 Flink 的新手,所以用一个明显的答案来问这个问题可能是个愚蠢的问题,但我目前不知道如何在这里实现我的目标。

编辑

按照大卫的回答,我尝试使用 Flink Table API 对数据流进行排序,使用 nc -lk 9999 进行文本套接字流。问题是在我关闭文本套接字流之前,没有任何东西打印到控制台。这是我写的scala代码-


package processor

import org.apache.flink.api.common.eventtime.{SerializableTimestampAssigner, WatermarkStrategy}
import org.apache.flink.api.common.functions.MapFunction
import org.apache.flink.api.common.state.{ValueState, ValueStateDescriptor}
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.KeyedProcessFunction
import org.apache.flink.streaming.api.scala._
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.table.api.{EnvironmentSettings, FieldExpression, WithOperations}
import org.apache.flink.util.Collector

import java.time.Duration


object AnnotateJob {

  val OUT_OF_ORDER_NESS = 100

  def main(args: Array[String]) {
    // set up the streaming execution environment
    val env = StreamExecutionEnvironment.getExecutionEnvironment

    val bSettings = EnvironmentSettings.newInstance().useBlinkPlanner().inStreamingMode().build()

    val tableEnv = StreamTableEnvironment.create(env, bSettings)

    env.setParallelism(1)

    // Obtain the input data by connecting to the socket. Here you want to connect to the local 9999 port.
    val text = env.socketTextStream("localhost", 9999)
    val objStream = text
      .filter( _.nonEmpty )
      .map(new ParticleMapFunction)

    val posStream = objStream
      .assignTimestampsAndWatermarks(
        WatermarkStrategy
          .forBoundedOutOfOrderness[ParticlePos](Duration.ofMillis(OUT_OF_ORDER_NESS))
          .withTimestampAssigner(new SerializableTimestampAssigner[ParticlePos] {
            override def extractTimestamp(t: ParticlePos, l: Long): Long = t.t
          })
      )

    val tablePos = tableEnv.fromDataStream(posStream, $"t".rowtime() as "et", $"t", $"name", $"x", $"y")
    tableEnv.createTemporaryView("pos", tablePos)
    val sorted = tableEnv.sqlQuery("SELECT t, name, x, y FROM pos ORDER BY et ASC")

    val sortedPosStream = tableEnv.toAppendStream[ParticlePos](sorted)

    // sortedPosStream.keyBy(pos => pos.name).process(new ValAnnotator)

    sortedPosStream.print()

    // execute program
    env.execute()
  }

  case class ParticlePos(t : Long, name : String, x : Double, y : Double) extends Serializable
  case class ParticlePosVal(t : Long, name : String, x : Double, y : Double,
                            var vx : Double = 0.0, var vy : Double = 0.0) extends Serializable

  class ParticleMapFunction extends MapFunction[String, ParticlePos] {
    override def map(t: String): ParticlePos = {
      val parts = t.split("\\W+")
      ParticlePos(parts(0).toLong, parts(1), parts(2).toDouble, parts(3).toDouble)
    }
  }

}

【问题讨论】:

    标签: apache-flink flink-streaming stream-processing


    【解决方案1】:

    一般来说,水印结合事件时间计时器是解决乱序事件流所带来问题的解决方案。涵盖Event Time and Watermarks 的官方 Flink 培训部分解释了这是如何工作的。

    在更高的层次上,有时使用 Flink 的 CEP 库或 Flink SQL 之类的东西会更容易,因为它们可以很容易地按时间对流进行排序,从而消除所有的乱序。例如,请参阅How to sort a stream by event time using Flink SQL 了解使用 Flink SQL 按事件时间对流进行排序的 Flink DataStream 程序示例。

    在您的情况下,一个相当简单的MATCH_RECOGNIZE 查询就可以满足您的需求。这可能看起来像这样,

    SELECT *
        FROM event
        MATCH_RECOGNIZE (
            PARTITION BY particleId
            ORDER BY ts
            MEASURES 
                b.ts, 
                b.particleId, 
                velocity(a, b)
            AFTER MATCH SKIP TO NEXT ROW
            PATTERN (a b)
            DEFINE
                a AS TRUE,
                b AS TRUE
        )
    

    velocity(a, b) 是一个用户定义的函数,用于计算速度,给定同一粒子的两个连续事件(a 和 b)。

    【讨论】:

    • 嘿大卫我关注了你链接的帖子,并尝试为我的用例实现解决方案,为我的数据流使用套接字,在我关闭之前没有任何东西打印到控制台套接字连接,我已经用我正在使用的代码编辑了我的问题,任何帮助将不胜感激
    • 这意味着您的作业看到的第一个水印是由套接字关闭生成的,它会在 MAX_WATERMARK 时间注入水印。可能有两件事之一是阻止创建较早的水印:要么没有足够大的时间戳的事件,要么作业没有运行足够长的时间(默认情况下,有界无序策略是每 200 毫秒调用一次新水印)。
    • 谢谢大卫,我没有任何时间戳足够大的事件,我现在可以看到输出
    【解决方案2】:

    在 Flink 中执行此操作的一种方法可能是使用 KeyedProcessFunction,即可以:

    • 处理信息流中的每个事件
    • 保持一些状态
    • 使用基于事件时间的计时器触发一些逻辑

    所以它会是这样的:

    • 您需要了解有关数据的某种“最大无序”。根据您的描述,我们假设例如 100 毫秒,这样在处理时间戳 1612103771212 的数据时,您决定认为您确定在 1612103771112 之前已收到所有数据。
    • 你的第一步是keyBy()你的流,通过粒子ID键控。这意味着您的 Flink 应用程序中的 next 运算符的逻辑现在可以用仅一个粒子的一系列事件来表示,并且每个粒子都以这种方式并行处理。

    类似这样的:

    yourStream.keyBy(...lookup p1 or p2 here...).process(new YourProcessFunction())
    
    • 在初始化ProcessFunctionYourProcessFunction 期间(即在open 方法期间),初始化一个ListState,您可以在其中安全地存储内容。
    • 在处理流中的元素时,在 processElement 方法中,只需将其添加到 listState 并在例如 100 毫秒内注册一个计时器触发器
    • onTimer() 方法触发时,比如在时间t,查看listState 中所有时间t - 100 的元素,如果您至少有两个元素,对它们进行排序,删除从状态中提取它们,应用您描述的速度计算和注释逻辑,并将结果发送到下游。

    您会发现an example in the official Flink training 在乘坐出租车期间使用这种逻辑,这与您的用例有很多相似之处。另请查看该 repo 的各种 Readme.md 文件以了解更多详细信息。

    【讨论】:

    • 已经在做keyBy(particleID)了,我将通过flink培训示例,看看它如何应用到我的案例中
    猜你喜欢
    • 1970-01-01
    • 2018-11-05
    • 1970-01-01
    • 2021-10-05
    • 1970-01-01
    • 2017-03-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多