I have 2 large files in both file i have the lat and long. i am trying to filter out all records that are within 500 meter.
I have follow the below step.
1.First I have joined with dataframe on the condition like abs(lat1-lat2) < 0.01 and abs(long1 - long2) < 0.01
val joindf = df1.join(broadcast(df2),abs(df1.lat1-df2.lat2) < 0.01 && abs(df1.long1-df2.long2) < 0.01) //will do the kind of crossjoin and return all rows within 5 km.
2.after that i have calculated the distance(basically call custum function that takes lat1,lat2,long1,long2 and its return distance in meter).
3.after that i have added the filter condition like distance <=500
But the above step works fine in small dataset but not working on Huge data. Its takes around 4 days in step 1.
Please help me out can we resolve this using geospatail. i have read the document but i am new in spark. please help me out.
The data frame sample like.
Dataframe 1
Tehsil District Type Code POP V_Lat V_Long
Tulsipur Balrampur NON Census Village 0 0 27.594705 82.334491
Tulsipur Balrampur NON Census Village 0 0 27.605287 82.34746
Tulsipur Balrampur NON Census Village 0 0 27.573511 82.336592
Tulsipur Balrampur NON Census Village 0 0 27.582564 82.355718
Tulsipur Balrampur NON Census Village 0 0 27.57687 82.322748
Tulsipur Balrampur NON Census Village 0 0 27.583982 82.344223
Tulsipur Balrampur NON Census Village 0 0 27.577273 82.330141
Tulsipur Balrampur NON Census Village 0 0 27.569862 82.326575
Tulsipur Balrampur Village 173435 2702 27.584897 82.353102
Tulsipur Balrampur Village 173434 2552 27.592387 82.330867
Tulsipur Balrampur Village 173436 3506 27.5734 82.340243
Tulsipur Balrampur Village 173431 1693 27.599005 82.345086
Haidergarh Bara Banki NON Census Village 0 0 26.579465 81.461515
dataframe 2
UE Purwa Dhanauti Haidergarh Bara Banki NON Census Village 0 0 26.568228 81.471936
UE Purwa Lachhmansingh Haidergarh Bara Banki NON Census Village 0 0 26.569711 81.478505