I find it pretty confusing using a broadcasted variable inside a UDF from an imported function. Say I make a broadcasted variable inside an imported function from the main file. It works if I have a UDF defined inside the function (second_func) but not outside (third_func).
- Why is this happening?
- Are UDFs advised being defined inside the function that calls it?
# test_utils.py
from pyspark.sql import types as T
from pyspark.sql import functions as F
@F.udf(T.StringType())
def do_smth_out():
return broadcasted.value["a"]
def second_func(spark, df):
@F.udf(T.StringType())
def do_smth_in():
return broadcasted.value["a"]
data = {"a": "c"}
sc = spark.sparkContext
broadcasted = sc.broadcast(data)
return df.withColumn("a", do_smth_in())
def third_func(spark, df):
data = {"a": "c"}
sc = spark.sparkContext
broadcasted = sc.broadcast(data)
return df.withColumn("a", do_smth_out())
# main.py
from pyspark.sql import types as T
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from test_utils import first_func, second_func
@F.udf(T.StringType())
def do_smth():
return broadcasted.value["a"]
if __name__ == "__main__":
spark = SparkSession \
.builder \
.getOrCreate()
sc = spark.sparkContext
columns = ["language","users_count"]
data = [("Java", "20000"), ("Python", "100000"), ("Scala", "3000")]
df = sc.parallelize(data).toDF(columns)
broadcasted = sc.broadcast({"a": "c"})
print("First trial")
df.withColumn("a", do_smth()).show()
# Works
print("Second trial")
second_func(spark, df).show()
# Works
print("Third trial")
third_func(spark, df).show()
# Doesn't work