Understanding Spark broadcasting

Viewed 62

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).

  1. Why is this happening?
  2. 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
    
0 Answers
Related