【问题标题】:flink - using dagger injections - not serializable?flink - 使用匕首注入 - 不可序列化?
【发布时间】:2015-12-06 14:26:23
【问题描述】:

我正在使用 Flink(最新通过 git)从 kafka 流式传输到 cassandra。为了简化单元测试,我通过 Dagger 添加了依赖注入。

ObjectGraph 似乎设置正确,但 Flink 将“内部对象”标记为“不可序列化”。如果我直接包含这些对象,它们会起作用 - 那么有什么区别?

有问题的类实现了 MapFunction@Inject 一个用于 cassandra 的模块和一个用于读取配置文件的模块。

有没有办法构建它,以便我可以使用后期绑定,还是 Flink 使这成为不可能?


编辑:

fwiw - 依赖注入(通过 dagger)和 RichMapFunction 不能共存。 Dagger 不允许您在定义中包含任何具有 extends 的对象。

进一步:

通过 Dagger Lazy 实例化的对象也不会序列化。

线程“主”org.apache.flink.api.common.InvalidProgramException 中的异常:对象 com.someapp.SaveMap@2e029d61 不可序列化
...
引起:java.io.NotSerializableException: dagger.internal.LazyBinding$1

【问题讨论】:

    标签: java serialization dagger apache-flink


    【解决方案1】:

    在深入探讨问题的细节之前,先了解一下 Apache Flink 中函数可序列化的背景:

    可序列化

    Apache Flink 使用 Java 序列化 (java.io.Serializable) 将函数对象(此处为 MapFunction)传送给并行执行它们的工作人员。因此,函数需要可序列化:函数可能不包含任何不可序列化的字段,即非原始类型(int、long、double、...)且未实现java.io.Serializable

    使用不可序列化构造的典型方法是延迟初始化它们。

    延迟初始化

    在 Flink 函数中使用不可序列化类型的一种方法是延迟初始化它们。当函数被序列化以交付时,保存这些类型的字段仍然是null,并且只有在函数被worker反序列化后才设置。

    • 在 Scala 中,您可以简单地使用惰性字段,例如 lazy val x = new NonSerializableType()NonSerializableType 类型实际上仅在第一次访问变量 x 时创建,该变量通常在工作程序上。因此,该类型可以是不可序列化的,因为x 在函数被序列化以传送给工作人员时为空。

    • 在 Java 中,您可以在函数的 open() 方法上初始化不可序列化的字段,如果您将其设为 Rich Function。丰富的函数(如RichMapFunction)是基本函数(此处为MapFunction)的扩展版本,可让您访问生命周期方法,如open()close()

    惰性依赖注入

    我对依赖注入不太熟悉,但 dagger 似乎也提供了类似惰性依赖的东西,这可能有助于作为一种解决方法,就像 Scala 中的惰性变量一样:

    new MapFunction<Long, Long>() {
    
      @Inject Lazy<MyDependency> dep;
    
      public Long map(Long value) {
        return dep.get().doSomething(value);
      }
    }
    

    【讨论】:

    • 可能更适合另一个问题,但 openclose 何时调用 RichFunction 运算符集?
    【解决方案2】:

    我遇到了类似的问题。有两种方法可以不反序列化您的依赖项。

    1. 使您的依赖关系静态化,但这并不总是可行的。它还会弄乱你的代码设计。

    2. 使用瞬态:通过将依赖项声明为瞬态,您是在说它们不是对象持久状态的一部分,也不应该成为序列化的一部分。

    public ClassA implements Serializable{
      //class A code here
    }
    
    public ClassB{
      //class B code here
    }
    
    public class MySinkFunction implements SinkFunction<MyData> {
      private ClassA mySerializableDependency;
      private transient ClassB nonSerializableDependency;
    }
    

    这在您使用外部库时特别有用,您无法更改其实现以使其可序列化。

    【讨论】:

    • 是的,但是如何初始化 nonSerializableDependency 呢?
    • 在 Flink 中,您可以从实现的丰富函数接口中使用 open() 方法。或者用一些初始值声明变量。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-08
    • 2015-11-04
    • 1970-01-01
    • 1970-01-01
    • 2017-06-03
    • 1970-01-01
    相关资源
    最近更新 更多