I am using spark 3.1.1 and joining two Dataframes of file size 8.6Gb and 25.2Mb respectively and not applying any filter. Spark is automatically using BroadcastHashJoin for this, although spark.sql.autoBroadcastJoinThreshold defaults to 10Mb.
How is 25.2Mb converted to 8.1Mb without applying any filter to be eligible for broadcast?
val df1 = spark.read
.option("header",true)
.csv("s3a://data/staging/received/data/spark/3/KernelVersionOutputFiles.csv")
.withColumn("Pid",substring(rand(),3,4).cast("bigint"))
val df2 = spark.read
.option("header",true)
.csv("s3a://data/staging/received/data/spark/3/ForumTopics.csv")
.withColumn("Cid",substring(rand(),3,4).cast("bigint"))
val df3 = df2.coalesce(1)
val joinDf = df1.join(df3, df1("Pid") === df3("Cid"))
val cnt = joinDf.count()
DAG looks like this:


