使用geopy处理160万推特地址效率低下,求提速方案
问题描述
我在项目中还原了160万条推特元数据(含创建日期、位置等),仅需关注美国地区的推特并按州生成统计数据。但多数位置信息错误或不规范,需要先规范化处理。
现有位置示例:
['한국어 강제 수용소 (DPRK)', 'Lagos, Nigeria', 'Kolkata, India', 'Who cares', 'Unknown', 'British Columbia, Canada', 'Bitcoin & Markets', 'White Plains, NY', 'Washington, DC']
我编写的筛选及规范化代码处理速度极慢(2.03it/s),处理全部数据需8-9天,寻求提速方案。当前测试代码如下:
from geopy import geocoders geolocator = geocoders.Nominatim(user_agent='myapplication') from tqdm.auto import tqdm tqdm.pandas() def get_adress(x): try: return geolocator.geocode(x).address except: return "" df_s = df.sample(1000) df_s["new_loc"] = df_s.user_location.progress_apply(get_adress) df_s["country"] = df_s.new_loc.apply(lambda x: x.split(",")[-1]) df_s = df_s[df_s.country.apply(lambda x: "United States" in x)] df_s = df_s[df_s.new_loc.apply(lambda x: len(x.split(","))) > 1]
提速方案
1. 前置过滤无效位置
先排除明显非美国、无意义的位置,大幅减少后续地理编码API的调用次数:
- 过滤含明确非美国国家名的字符串(如Nigeria、India、Canada等)
- 过滤无意义内容:如"Who cares"、"Unknown"、乱码文本、与地理位置无关的字符串(如"Bitcoin & Markets")
示例代码:
# 定义无效位置关键词集合 invalid_keywords = {"Unknown", "Who cares", "Bitcoin & Markets", "Nigeria", "India", "Canada", "DPRK"} def is_valid_location(loc): if not isinstance(loc, str): return False # 排除含无效关键词的内容 for kw in invalid_keywords: if kw in loc: return False # 排除乱码(简单判断:非ASCII字符占比过高) ascii_ratio = sum(1 for c in loc if ord(c) < 128) / len(loc) return ascii_ratio >= 0.5 # 先过滤无效位置,缩减待处理数据集 df_filtered = df[df["user_location"].apply(is_valid_location)]
2. 缓存+速率限制优化API调用
- 缓存机制:用
functools.lru_cache缓存已查询过的位置结果,避免重复请求相同内容 - 速率限制:遵循Nominatim官方规则添加请求间隔,避免被限流封禁
示例代码:
from functools import lru_cache from geopy.geocoders import Nominatim from geopy.extra.rate_limiter import RateLimiter geolocator = Nominatim(user_agent='myapplication') # 设置速率限制,避免触发API限流 geocode = RateLimiter(geolocator.geocode, min_delay_seconds=1) # 缓存查询结果,最大缓存10万条 @lru_cache(maxsize=100000) def get_address_cached(loc): try: # 限制搜索范围为美国,缩小API检索范围 res = geolocator.geocode(loc, country_codes='us') return res.address if res else "" except: return "" # 应用到过滤后的数据集 df_filtered["new_loc"] = df_filtered["user_location"].progress_apply(get_address_cached)
3. 并行处理提升效率
利用多进程或swifter库实现并行化处理,充分利用CPU资源:
- 用
swifter自动选择最优处理方式:
import swifter df_filtered["new_loc"] = df_filtered["user_location"].swifter.apply(get_address_cached)
- 手动多进程处理(注意:多进程下缓存需用共享内存或外部存储,如Redis):
from multiprocessing import Pool def process_batch(batch): return [get_address_cached(loc) for loc in batch] # 分批次处理,避免内存溢出 batch_size = 1000 batches = [df_filtered["user_location"].iloc[i:i+batch_size] for i in range(0, len(df_filtered), batch_size)] with Pool(processes=4) as pool: results = pool.map(process_batch, batches) # 合并结果 df_filtered["new_loc"] = [item for sublist in results for item in sublist]
4. 替换为商业地理编码服务(可选)
如果预算允许,改用Google Maps Geocoding API、Mapbox Geocoding API等商业服务,这类服务响应速度更快、并发量更高,适合大规模数据处理,能大幅缩短处理时间。
内容的提问来源于stack exchange,提问作者lifrah
相关产品推荐
相关产品推荐

