如何在Spark Scala UDF中获取调用查询及Spark配置?
在Spark/Scala(Databricks)的UDF中获取调用查询与Spark配置
问题背景
几年前在Hadoop环境编写Hive UDF时,可通过以下方式获取调用查询:
public class MyHiveUDF extends GenericUDF { @Override public Object evaluate(DeferredObject[] arguments) throws CryptoException { // 业务逻辑实现 } @Override public void configure(MapredContext context) { super.configure(context); Configuration conf = context.getJobConf(); String myQuery = conf.get("hive.query.string",""); } }
当前在Spark 3.3.0、Scala 2.12、Databricks 11.x环境下,希望在Scala UDF中实现类似功能,期望代码大致如下:
def myCustomUdf: UserDefinedFunction = udf ( (colValue : String) => { val myQuery = ...// 此处获取查询字符串 stuffWithInput(colValue) })
核心问题:
- 是否能在Spark UDF中获取调用查询?若可行,如何实现?
- 是否能在UDF代码中获取任意Spark配置?
解决方案
1. 获取调用查询字符串
Spark不会自动将当前执行的查询字符串传入UDF的执行上下文,直接在UDF内部实时获取是不可行的,但可以通过以下间接方式实现:
- 提前传递查询字符串:在执行UDF前,先获取当前查询内容(笔记本环境可通过Databricks的
dbutils.notebook.getContext()关联作业日志提取查询;提交的作业可通过Databricks作业API获取对应查询文本),再将查询字符串作为参数传入UDF。
示例代码:// 假设已通过对应方式获取到查询字符串 val queryString = "SELECT myCustomUdf(col) FROM my_table" def myCustomUdf: UserDefinedFunction = udf ( (colValue: String, query: String) => { // 结合查询字符串处理业务逻辑 stuffWithInput(colValue, query) }) // 使用时传入查询字符串 df.select(myCustomUdf($"col", lit(queryString))) - 事后通过事件日志关联:如果是离线作业,可通过Spark事件日志(Event Log)关联作业ID与对应的查询内容,但这种方式无法在UDF运行时实时获取,仅能用于事后分析。
2. 获取Spark配置
UDF运行在Executor节点上,Driver端的SparkConf无法直接序列化传递到Executor,因此不能直接在UDF中访问完整的Spark配置,但可以通过以下方式实现配置获取:
- 提前传递目标配置项:在Driver端获取需要的配置值,通过UDF参数传入,或者将配置值广播(
Broadcast)后在UDF中使用。
广播方式示例:// 在Driver端获取配置并广播 val targetConfig = spark.sparkContext.getConf.get("spark.target.config.key", "default_value") val broadcastedConfig = spark.sparkContext.broadcast(targetConfig) def myCustomUdf: UserDefinedFunction = udf ( (colValue: String) => { // 在UDF中使用广播的配置值 val configVal = broadcastedConfig.value stuffWithInput(colValue, configVal) }) - 访问Executor系统属性:部分Spark配置会同步为Executor的系统属性,可在UDF中通过
System.getProperty("配置键名")获取,但仅适用于少数会同步的配置项,无法覆盖所有Spark配置。
内容的提问来源于stack exchange,提问作者Nitznutz
相关产品推荐
相关产品推荐

