【问题标题】:Case class serialazation in flinkflink中的案例类序列化
【发布时间】: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


【解决方案1】:

我的程序的定义使用了一个类

class Job(conf: AdfinConfig)(implicit env: ExecutionEnvironment)
        extends DspJob(conf){
    ...
    case class TestTarget(tacticId: String, partnerId:Long)
    campaignPartners.join(partnerInput).where(1).equalTo("id") {
    ...
}

因为它是一个内部类,它没有被自动序列化

如果您将类切换为不是内部类,那么一切正常

case class TestTarget(tacticId: String, partnerId:Long)
class Job(conf: AdfinConfig)(implicit env: ExecutionEnvironment)
        extends DspJob(conf){
    ...
    words.join( ....) 
    ...
}

【讨论】:

  • 这个问题也可以用 Kryo / RocksDB 序列化。当我尝试通过 Kryo 将内部 case class 序列化到 RocksDB 时,出现以下错误:com.esotericsoftware.kryo.KryoException: java.lang.ClassCastException: scala.Tuple2 cannot be cast to org.joda.time.DateTime 出于某种原因,Kryo 将案例类解释为 Tuple2 对象,并使错误变得毫无意义。通过将我所有的内部案例类移动为非内部案例类来解决问题。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-03-29
  • 1970-01-01
  • 2013-08-12
  • 2011-05-03
  • 2011-12-06
相关资源
最近更新 更多