gchandra
Databricks Employee
Databricks Employee

df1 has 1500 rows and df2 has 9 million rows, try broadcast join on df1

df= spark.sql("""SELECT /*+ BROADCAST(df1) */
df1.*, df2.* FROM df1 JOIN df2  ON ST_Intersects(df1.geometry, df2.geometry) """)

 



~