I have PySpark dataframes with couple of columns, on of them being gps location (in WKT format). What is the easiest way to pick only rows that are inside some polygon? Does it scale when there are ~1B rows?
I'm using Azure Databricks and if the solution exists in Python, that would be even better, but Scala and SQl are also fine.
Edit: Alex Ott's answer - Mosaic - works and I find it easy to use.