【问题标题】:How do I pass Spark broadcast variable to a UDF in Java?如何将 Spark 广播变量传递给 Java 中的 UDF?
【发布时间】:2020-02-03 10:06:53
【问题描述】:

我有一个 Properties 变量,我通过调用 REST 服务来填充它。我需要将此列表传递给我的 UDF。我认为广播变量可以很好地满足我的目的,因为属性列表可能很长。所以我在我的主类中创建了一个广播变量:

Properties kp = getApplicationProperties(rootPath);
Broadcast<Properties> brVar = sc.broadcast(kp);

但是,我将在我的 UDF 中发送什么变量?我尝试使用 org.apache.spark.sql.functions.lit 发送文字值,但这导致我的 UDF 从未被调用:

spark.sqlContext().udf().registerJava("MongoInsert", MongoInsert.class.getName(), DataTypes.StringType);
persondatasetwithResid.select(callUDF("MongoInsert", lit(rootPath).cast(DataTypes.StringType),
                            col("value").cast(DataTypes.StringType), col("resourceId").cast(DataTypes.StringType),
                            lit(brVar))).show();


public class MongoInsert implements UDF4<String, String, String, Broadcast<Properties>, String> {
  public String call(String rootPath, String jsonstring, String resourceId, Broadcast<Properties> brVar)
        throws Exception {

    Properties kp = brVar.value();
  }
}

我找不到任何可以解释如何在 Java 中将广播变量传递给 UDF 的资源。请帮帮我。

编辑

我在网络上得到指示,只有列类型和文字字符串类型可以传递到 UDF。是这样吗?可以不传入其他变量,如地图、数组等吗?请记住,我说的是 Spark 和 Java。在 Scala 中,似乎有更多的灵活性。

编辑 2

我也收到了很多关于 typedLit 的文献,它实际上可以让我们使用 Maps 和 Arrays,例如。以下问题:

How to add a Map column to Spark dataset?

这是否意味着将变量包装在 lit() 中不会达到我的目的?我尝试用 typedLit() 包装一个简单的 Map 变量,但它给了我一个编译错误,说

"The method typedLit(T, TypeTags.TypeTag<T>) in the type functions is not applicable for the arguments (Map<String,String>)"

可以预见,网络上有大量关于如何在 Scala 中使用 typedLit 的资源,但在 Spark 和 Java 中几乎没有。

编辑 3

我发现了另一个问题:

how to set and get static variables from spark?

这个也可以为我提供答案,因为传递给所有类的静态变量可以满足我的目的。答案再次指定了广播变量,但也指定了闭包。再一次,在 Java 中没有使用这种闭包的例子,甚至在 Spark 官方文档中也没有!如果有人可以告诉我如何在 Java 中创建一个闭包并使用它将一个变量传递给 UDF,这将对我有很大帮助。

【问题讨论】:

  • scala 语法对你有用吗?相同的概念
  • 我认为你这样做的方式是正确的。您必须始终将广播包装器传递给 udf 并在函数中访问其值。
  • @user238607,即使方法是正确的,实现也肯定会被某些东西关闭,因为没有调用 UDF。

标签: java apache-spark


【解决方案1】:

这是一个在类级别访问广播变量而不是被传入的示例。

import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.broadcast.Broadcast;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.api.java.UDF1;
import org.apache.spark.sql.types.DataTypes;

public class Test {

    Broadcast<String> broadcastVar = null;

    public UDF1 myudf =  new UDF1<String,String>(){
        @Override
        public String call(String x) {
            return broadcastVar.getValue();
        }
    };

    public void setBroadcastVar(Broadcast<String> broadcastVar){
        this.broadcastVar = broadcastVar;
    }


    public static void main(String[] args) {

        SparkSession spark = SparkSession.builder().master("local").appName("test").getOrCreate();
        JavaSparkContext js = new JavaSparkContext(spark.sparkContext());

        Test t = new Test();
        t.setBroadcastVar(js.broadcast("hellow world"));

        spark.udf().register("myudf",t.myudf,DataTypes.StringType);

       //use the udf

    }

}

【讨论】:

  • 感谢您的回答,萨利姆。但是,我的 UDF 使用它们的调用方法声明为它们自己的单独类。这就是为什么我需要以某种方式将广播变量传递给它们。
  • 嗨萨利姆。我搜索了如何将静态变量发送到 UDF,发现了这个问题:stackoverflow.com/questions/29685330/…,一个答案说驱动程序中定义的静态变量在执行程序中不可见,并且广播变量或闭包是唯一的方法。像往常一样,我也没有得到太多带有 Java 闭包的 Spark 示例。
  • 你试过这种方式吗?在类上定义一个静态变量,用广播值填充该变量,在 udf 类中引用该变量。
  • 嗨,Salim,很抱歉回复晚了,但与此同时我被拉到了其他任务中。我复制了你的类,由于这里的 UDF 被定义为匿名内部类,它能够访问外部作用域的广播变量。但是,我创建了与单独的类相同的 UDF,并从那里调用了变量,然后立即引发了编译错误。我忽略了它并尝试运行该程序,但显示 Spark 异常堆栈跟踪说未解决的编译问题:broadcastVar 无法解决
  • 您在 udf 类中定义广播变量并使用结算器方法设置它。复用类的实例多次调用udf
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-06
  • 1970-01-01
  • 1970-01-01
  • 2018-02-02
  • 1970-01-01
  • 2018-05-10
相关资源
最近更新 更多