Spark Scala:如何基于另一DataFrame为用户表新增大小写不敏感匹配的国家列
解决方案
一、Pandas 实现方式
步骤1:预处理城市-国家映射表
先把映射表的城市名统一转为小写(或大写),生成查询字典,提升匹配效率:
import pandas as pd # 示例映射DataFrame city_country_df = pd.DataFrame({ 'city': ['Washington', 'London', 'Paris'], 'country': ['USA', 'UK', 'France'] }) # 生成小写城市到国家的映射字典 city_country_map = city_country_df.set_index(city_country_df['city'].str.lower())['country'].to_dict()
步骤2:为用户表新增country列
方式1:调用自定义函数foo
def foo(city): # 将输入城市转小写后匹配字典,未匹配到返回None或自定义默认值 return city_country_map.get(city.lower(), None) # 示例用户DataFrame user_df = pd.DataFrame({ 'user_id': [1, 2, 3], 'city': ['washington', 'London', 'Berlin'] }) # 新增country列 user_df['country'] = user_df['city'].apply(foo)
方式2:直接用map(更简洁高效)
无需自定义函数,直接转小写后映射:
user_df['country'] = user_df['city'].str.lower().map(city_country_map)
二、PySpark 实现方式
Spark中优先用关联查询(UDF性能较差),实现大小写不敏感匹配:
步骤1:预处理映射表,新增小写城市列
from pyspark.sql import SparkSession from pyspark.sql.functions import lower spark = SparkSession.builder.appName("city-country-match").getOrCreate() # 示例映射DataFrame city_country_df = spark.createDataFrame([ ("Washington", "USA"), ("London", "UK"), ("Paris", "France") ], ["city", "country"]) # 新增小写城市列用于关联 city_country_df = city_country_df.withColumn("city_lower", lower("city"))
步骤2:关联用户表获取国家信息
# 示例用户DataFrame user_df = spark.createDataFrame([ (1, "washington"), (2, "London"), (3, "Berlin") ], ["user_id", "city"]) # 用户表转小写后关联,清理冗余列 user_df = user_df.withColumn("city_lower", lower("city")) \ .join(city_country_df, on="city_lower", how="left") \ .drop("city_lower", city_country_df["city"])
若必须使用自定义UDF(不推荐)
# 将映射表转为字典并广播到集群,避免重复加载 city_country_map = city_country_df.select("city_lower", "country").toPandas().set_index("city_lower")["country"].to_dict() broadcast_map = spark.sparkContext.broadcast(city_country_map) # 定义并注册UDF from pyspark.sql.functions import udf from pyspark.sql.types import StringType def foo(city): return broadcast_map.value.get(city.lower(), None) foo_udf = udf(foo, StringType()) # 新增country列 user_df = user_df.withColumn("country", foo_udf("city"))
内容的提问来源于stack exchange,提问作者Guelque Hadrien
相关产品推荐
相关产品推荐

