【问题标题】:Beam DoFn static variable shared across JVM跨 JVM 共享的 Beam DoFn 静态变量
【发布时间】:2020-04-10 20:07:30
【问题描述】:

所以,我试图弄清楚 Beam DoFn 中静态变量的行为, 它是否在线程之间共享(在同一个 JVM 中)?

基本上试图从编程指南中理解以下内容:

4.3.2。线程兼容性
…请注意,函数对象中的静态成员不会传递给工作实例,并且多个 可以从不同的线程访问您的函数的实例。

https://beam.apache.org/documentation/programming-guide/#requirements-for-writing-user-code-for-beam-transforms

现在看来,下面的静态对象“counter”在worker(Flink引擎)中被初始化、序列化和应用了,它是否与上面的语句一致?

如果工作线程属于不同的进程/JVM,显然不会被共享。但是如果掉到同一个 JVM 会不会共享“计数器”?

public class myTransform extends DoFn<KV<String >,String> implements Serializable {
    private static AtomicLong counter = new AtomicLong(0);
         ...
         @ProcessElement
         public void processElement(ProcessContext c) {
             ...
             counter.incrementAndGet();
         }
}

谢谢

【问题讨论】:

    标签: apache-flink apache-beam


    【解决方案1】:

    我认为初始化部分是指例如在DoFn 的构造函数或其他东西中设置一些值。您的代码将被初始化,因为 Worker 必须加载 myTransform 类。

    如果它们碰巧在同一个 JVM 中运行,那么是的,这将是共享的。 Beam 人试图传达的是,无论如何你都不应该以此为基础,并且操作符的并行实例可能会在任何节点上执行。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2014-08-09
      • 1970-01-01
      • 2012-05-19
      • 2022-11-18
      • 2011-01-31
      • 2011-03-24
      • 1970-01-01
      相关资源
      最近更新 更多