pythonno_udf = (flights
.filter(F.col("year") == 2024)
.withColumn("band", band)
.join(carriers, "carrier")
.groupBy("carrier_name", "band").count())
broadcast_too = (flights
.filter(F.col("year") == 2024)
.withColumn("band", band)
.join(broadcast(carriers), "carrier")
.groupBy("carrier_name", "band").count())
change_2 = measure("2. built-in functions for the band", lambda: noop(no_udf))
change_3 = measure("3. broadcast the carriers", lambda: noop(broadcast_too))
print("same answer :", sorted(slow.collect()) == sorted(broadcast_too.collect()))
text = plan_text(broadcast_too)
print("PartitionFilters on year :", "PartitionFilters: [isnotnull(year" in text)
print("BatchEvalPython :", "BatchEvalPython" in text)
print("join, slow query :", join_used(slow))
print("join, after the changes :", join_used(broadcast_too))
# The built-in band removes the Python node. The hint makes the carriers the side that is copied :
# without it Spark shuffles both sides, or on a small copy of the flights may broadcast the flights.