You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 23:30:43