PySpark UDF中是否仍需广播模型mdl?UDF会自动处理吗?
关于PySpark UDF中模型广播的疑问
问题描述
我有一个模型mdl,通过代码mdl = spark.sparkContext.broadcast(mdl).value进行广播。我编写了一个PySpark UDF来调用该模型:
@udf(returnType=StringType()) def func(row): result = mdl.predict(row.features) # 其余代码
请问使用该UDF时,我是否仍需对模型mdl进行广播?还是UDF会自动处理这一操作?
解答
你当前的写法其实是无效的广播操作——spark.sparkContext.broadcast(mdl).value会先将模型广播,随即直接取出广播对象的本地值重新赋值给mdl,相当于mdl还是原来的本地对象,完全没起到广播的作用。
另外要明确:PySpark的UDF不会自动处理外部变量的广播。如果不对模型做正确的广播配置,每个Executor在执行UDF时都会单独拉取一份模型副本,既浪费内存资源,又会拖慢任务执行效率。
正确的做法
- 先保留广播对象,不要直接取
.value:
broadcast_mdl = spark.sparkContext.broadcast(mdl)
- 在UDF内部通过广播对象的
.value获取模型实例:
@udf(returnType=StringType()) def func(row): result = broadcast_mdl.value.predict(row.features) # 其余代码
这样每个Executor只会持有一份模型副本,避免了重复传输和内存冗余。
内容的提问来源于stack exchange,提问作者ARIJIT SINGH
相关产品推荐
相关产品推荐

