【问题标题】:Filter Flink tuples过滤 Flink 元组
【发布时间】:2018-01-15 13:54:33
【问题描述】:

我正在使用 Flink 在 Scala 中编写流处理程序。我有一个数据流,我首先将其映射到包含 json4s JValues 的元组。现在我想根据这些 JValue 过滤这些元组。我认为这很简单,但我找不到任何关于如何按列过滤 Flink 元组的好例子。 有谁知道如何做到这一点? 谢谢

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    您可以简单地映射到case classes 并过滤掉不需要的东西,而不是映射到元组:

    // StreamingJob.scala
    
    ...
    
    val filteredEvents = content
          .map(x => Event.toCaseClass(x))
          .filter(x => x.value == true)
    
    ...
    
    // Event.scala
    
    case class Event(
                      id: String,
                      value: Int,
                    )
    object Event {
      implicit val formats = DefaultFormats
    
      def toCaseClass(str: String) =
        parse(str).extract[Event]
    }
    

    【讨论】:

      【解决方案2】:

      这个问题对我来说似乎有点太不确定了,但也许这不起作用?

      // stream contains stuff like these in a flink tuple 
      //(custom deserializer of array to tuple2???)
      val jsonExample = """["foo", "bar"]"""
      
      val stream: DataStream[Tuple2[JString, JString]] = ???
      val filteredStream = stream.filter(x => x.getField(0).extract[String] == "foo")
      

      如果您正在编写 scala,最好不要使用 flink 元组。选择案例类或至少是 scala 元组?

      【讨论】:

      • 我在一项研究作业中提出了这个问题,其中指出我必须使用 flink 元组。也许这不是最好的方法,但任务迫使我使用它。
      猜你喜欢
      • 2018-12-09
      • 2016-09-11
      • 1970-01-01
      • 2016-01-10
      • 2013-05-03
      • 1970-01-01
      • 2014-09-11
      • 2017-09-29
      相关资源
      最近更新 更多