Spark UDF结合Geopy逆地理编码处理10条数据陷入无限执行
问题原因及解决方案
核心问题分析
你的UDF在小样本测试正常,但处理真实数据(哪怕仅10行)时陷入无限执行,主要有以下几个关键原因:
- 重复创建Nominatim客户端:UDF每处理一行就初始化一个
Nominatim(user_agent="geoap")实例,频繁创建连接会触发Nominatim的访问限制,甚至被服务器临时封禁请求,导致任务无限等待。 - 无异常处理机制:真实数据中可能存在经纬度为null、无效坐标的情况,调用
location.raw['address']时会抛出未捕获的异常,Spark任务会因异常陷入重试或阻塞状态。 - Spark并行请求触发限流:Spark会并行执行UDF,即使只有10行,也可能同时发起多个请求。Nominatim对请求频率有严格限制(官方建议每秒最多1次请求),大量并发请求会被服务器限流,导致请求超时、重试,表现为任务"无限执行"。
修复后的代码示例
1. 复用Nominatim客户端(通过广播变量)
将客户端初始化移到UDF外部,用Spark广播变量传递,避免重复创建实例:
from geopy.geocoders import Nominatim import pyspark.sql.functions as F # 全局初始化一次客户端,广播到所有节点 geolocator = Nominatim(user_agent="geoap") broadcast_geolocator = spark.sparkContext.broadcast(geolocator) @F.udf("string") def city_state_country(lat,lng): # 从广播变量获取复用的客户端 locator = broadcast_geolocator.value # 处理null值 if not lat or not lng: return None coord = f"{lat},{lng}" try: # 设置超时时间,避免请求挂起 location = locator.reverse(coord, exactly_one=True, timeout=10) if not location: return "" address = location.get('address', {}) country = address.get('country', '') return country except Exception as e: # 捕获所有异常,返回标记值便于排查 return f"ERROR: {str(e)}"
2. 预处理数据,过滤无效值
在调用UDF前先过滤掉经纬度为空的行,减少无效请求:
# 过滤经纬度非空的行 df_open_app3 = df_open_app2.select("starting_lng","starting_lat")\ .filter(F.col("starting_lat").isNotNull() & F.col("starting_lng").isNotNull())\ .limit(10) df_open_app4 = df_open_app3.withColumn('con', city_state_country("starting_lat","starting_lng"))
3. 控制请求频率(可选)
如果仍遇到限流问题,可在UDF中添加延迟,符合Nominatim的请求频率要求:
import time @F.udf("string") def city_state_country(lat,lng): locator = broadcast_geolocator.value if not lat or not lng: return None coord = f"{lat},{lng}" try: time.sleep(1) # 控制每秒最多1次请求 location = locator.reverse(coord, exactly_one=True, timeout=10) if not location: return "" address = location.get('address', {}) country = address.get('country', '') return country except Exception as e: return f"ERROR: {str(e)}"
额外注意事项
- Nominatim是免费服务,不适合处理400万行级别的大规模数据,建议使用离线地理编码库(如
geopandas结合离线shapefile),或付费的地理编码API。 - 处理大规模数据时,先对经纬度去重,减少重复请求,能大幅提升效率。
内容的提问来源于stack exchange,提问作者BryC
相关产品推荐
相关产品推荐

