如何用PySpark为每个客户邮编匹配最近的仓库邮编?
用PySpark实现客户邮编匹配最近仓库的方案
核心思路
要实现这个需求,关键是先把邮编转换为经纬度,再通过Haversine公式计算球面距离,最后为每个客户筛选出距离最近的仓库。由于4万条客户数据×300个仓库=1200万条计算量,这个规模在PySpark的处理能力范围内,无需额外复杂优化。
步骤实现
1. 准备邮编-经纬度映射表
首先需要一个包含美国邮编对应经纬度的数据集(可从公开数据源获取后导入Spark),假设表结构为:
zip_code STRING, latitude DOUBLE, longitude DOUBLE
2. 加载并关联原始表与经纬度
先加载CUSTOMER_ORDERS和warehouse_loc表,再分别关联邮编经纬度表,得到带经纬度的客户和仓库数据:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, radians, sin, cos, sqrt, asin, row_number from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("NearestWarehouse").getOrCreate() # 加载原始表 customer_orders = spark.read.table("CUSTOMER_ORDERS") warehouse_loc = spark.read.table("warehouse_loc") # 加载邮编经纬度映射表(假设已导入Spark) zip_lookup = spark.read.table("zip_lookup") # 关联客户表与经纬度 customer_with_geo = customer_orders.join( zip_lookup, customer_orders.CUST_POSTAL_CD == zip_lookup.zip_code, "left" ).select( col("GEO"), col("CUST_POSTAL_CD").alias("customer_zip"), col("UNITS"), col("latitude").alias("cust_lat"), col("longitude").alias("cust_lon") ) # 关联仓库表与经纬度 warehouse_with_geo = warehouse_loc.join( zip_lookup, warehouse_loc.WH_ZIP == zip_lookup.zip_code, "left" ).select( col("WH_ID"), col("WH_ZIP").alias("warehouse_zip"), col("WH_TYPE"), col("latitude").alias("wh_lat"), col("longitude").alias("wh_lon") )
3. 计算所有客户-仓库对的距离
通过笛卡尔积关联客户和仓库数据,用Haversine公式计算球面距离(单位:公里):
# 定义Haversine距离计算函数 def haversine_distance(lat1, lon1, lat2, lon2): # 转换为弧度 lat1_rad = radians(lat1) lon1_rad = radians(lon1) lat2_rad = radians(lat2) lon2_rad = radians(lon2) dlat = lat2_rad - lat1_rad dlon = lon2_rad - lon1_rad a = sin(dlat/2)**2 + cos(lat1_rad) * cos(lat2_rad) * sin(dlon/2)**2 c = 2 * asin(sqrt(a)) # 地球半径,单位公里 r = 6371 return c * r # 注册为UDF haversine_udf = spark.udf.register("haversine", haversine_distance) # 笛卡尔积关联并计算距离 customer_warehouse_dist = customer_with_geo.crossJoin(warehouse_with_geo).withColumn( "distance_km", haversine_udf(col("cust_lat"), col("cust_lon"), col("wh_lat"), col("wh_lon")) )
4. 筛选每个客户的最近仓库
使用窗口函数按客户邮编分组,按距离升序排序,取每组的第一条数据:
# 定义窗口:按客户邮编分组,按距离升序排序 window_spec = Window.partitionBy("customer_zip").orderBy(col("distance_km").asc()) # 筛选最近仓库 nearest_warehouse = customer_warehouse_dist.withColumn( "row_num", row_number().over(window_spec) ).filter(col("row_num") == 1).select( col("GEO"), col("customer_zip").alias("CUST_POSTAL_CD"), col("UNITS"), col("WH_ID"), col("warehouse_zip").alias("WH_ZIP"), col("WH_TYPE"), col("distance_km").alias("nearest_distance_km") ) # 查看结果 nearest_warehouse.show()
注意事项
- 确保邮编-经纬度映射表覆盖所有客户和仓库的邮编,避免关联后出现空值;若有空值,可根据业务需求处理(比如过滤或标记)。
- Haversine公式的单位可调整:如果需要英里,将地球半径改为3956即可。
- 若后续数据量大幅增长(比如客户百万级),可以考虑先按区域粗筛再计算距离,减少笛卡尔积的计算量,但当前4万×300的规模无需此优化。
内容的提问来源于stack exchange,提问作者1ksj8jdnu36flksf
相关产品推荐
相关产品推荐

