【发布时间】:2018-01-15 13:54:33
【问题描述】:
我正在使用 Flink 在 Scala 中编写流处理程序。我有一个数据流,我首先将其映射到包含 json4s JValues 的元组。现在我想根据这些 JValue 过滤这些元组。我认为这很简单,但我找不到任何关于如何按列过滤 Flink 元组的好例子。 有谁知道如何做到这一点? 谢谢
【问题讨论】:
标签: scala apache-flink flink-streaming
我正在使用 Flink 在 Scala 中编写流处理程序。我有一个数据流,我首先将其映射到包含 json4s JValues 的元组。现在我想根据这些 JValue 过滤这些元组。我认为这很简单,但我找不到任何关于如何按列过滤 Flink 元组的好例子。 有谁知道如何做到这一点? 谢谢
【问题讨论】:
标签: scala apache-flink flink-streaming
您可以简单地映射到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]
}
【讨论】:
这个问题对我来说似乎有点太不确定了,但也许这不起作用?
// 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 元组?
【讨论】: