【发布时间】:2016-07-02 17:58:00
【问题描述】:
我正在尝试使用 Scala 的案例类构建数据集(我想在元组上使用案例类,因为我想按名称连接字段)。
这是我正在处理的连接的一次迭代:
case class TestTarget(tacticId: String, partnerId:Long)
campaignPartners.join(partnerInput).where(1).equalTo("id") {
(target, partnerInfo, out: Collector[TestTarget]) => {
partnerInfo.partner_pricing match {
case Some(pricing) =>
out.collect(TestTarget(target._1, partnerInfo.partner_id))
case None => ()
}
}
}
显然这会引发错误:
org.apache.flink.api.common.InvalidProgramException: 没有任务 可序列化在 org.apache.flink.api.scala.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:179) 在 org.apache.flink.api.scala.ClosureCleaner$.clean(ClosureCleaner.scala:171) 在 org.apache.flink.api.scala.DataSet.clean(DataSet.scala:121) 在 org.apache.flink.api.scala.JoinDataSet$$anon$2.(joinDataSet.scala:108) 在 org.apache.flink.api.scala.JoinDataSet.apply(joinDataSet.scala:107) 在 com.adfin.dataimport.vendors.dbm.Job.calculateVendorFees(Job.scala:84)
我已经看到文档here 指出我需要为该类实现可序列化。据我所知,在新版本的 Scala 中没有自动序列化案例类的好方法。 (我研究了手动序列化,但我认为我需要对链接做一些额外的工作才能使它工作)。
编辑: 根据 Till Rohrmann 的建议,我尝试使用一个小案例重现此错误。这是我用来尝试重现错误的方法。此示例有效,但我未能重现该错误。我还尝试将 Option 案例放在任何地方,但这也会导致工作失败。
val text = env.fromElements("To be, or not to be,--that is the question:--")
val words = text.flatMap { _.toLowerCase.split("\\W+") }.map(x => (1,x))
val nums = env.fromElements(List(1,2,3,4,5)).flatMap(x => x).map(x => First(1,x))
val counts = words.join(nums).where(0).equalTo("a") {
(a, b, out: Collector[TestTarget]) => {
b.b match {
case 2 => ()
case _ => out.collect(TestTarget(a._2, b.b))
}
}
}
【问题讨论】:
-
您能否提供一个完整的示例来重现您的问题?
-
您还需要什么?调用此函数后,我使用 writeAsText 将其输出。 campaignPartners 和 PartnerInfo 是 DataSets 你想要它们的类型签名吗?
-
我用 Flink
1.1-SNAPSHOT测试了你的例子,它工作得很好。如果错误仍然存在,您能否发布一个完整的示例,包括类型First以及您定义它们的位置(可以简单地复制和粘贴)。理想情况下,您只需发布 Scala 文件。 -
对不起,小例子在我无法重现错误中起作用
-
哦,对不起,我看错了。你的初始计划呢?
标签: scala serialization apache-flink