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

如何用PySpark提取JSON值并拼接成忽略空值的地址字符串?

PySpark 实现地址字段空值忽略与拼接

需求说明

已通过from_json解析得到包含id和addresses数组的DataFrame,其中addresses的每个元素是包含street、city、postalCode、state、country字段的结构体。需要将每个地址忽略空值字段后拼接为字符串,不同地址间用分号分隔。

实现方案

以下提供两种可行的实现方式,可根据数据量和复杂度选择:


方式一:使用内置函数(高效推荐)

利用PySpark内置的数组、Map处理函数完成,无需自定义UDF,性能更优:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, ArrayType
import pyspark.sql.functions as F

# 初始化Spark会话
spark = SparkSession.builder.appName("AddressFormatting").getOrCreate()

# 定义地址结构体Schema(与解析后的Schema一致)
address_schema = StructType([
    StructField("street", StringType(), True),
    StructField("city", StringType(), True),
    StructField("postalCode", StringType(), True),
    StructField("state", StringType(), True),
    StructField("country", StringType(), True)
])

# 模拟解析后的测试DataFrame
data = [
    ("1", [
        {"street": "123 Main St", "city": "New York", "postalCode": "10001", "state": "NY", "country": "USA"},
        {"street": None, "city": "Los Angeles", "postalCode": "90001", "state": None, "country": "USA"}
    ]),
    ("2", [
        {"street": "456 Oak Ave", "city": None, "postalCode": None, "state": "CA", "country": "USA"}
    ])
]
df = spark.createDataFrame(data, schema=StructType([
    StructField("id", StringType()),
    StructField("addresses", ArrayType(address_schema))
]))

# 核心处理逻辑
processed_df = df.withColumn(
    "formatted_addresses",
    # 用分号分隔多个地址
    F.array_join(
        # 遍历addresses数组,处理每个地址
        F.transform(
            F.col("addresses"),
            lambda addr: 
                # 拼接单个地址的非空字段(用逗号分隔)
                F.concat_ws(", ",
                    # 将过滤后的键值对转为"key:value"格式
                    F.transform(
                        # 过滤掉值为空的键值对
                        F.filter(
                            # 将地址结构体转为Map类型(键为字段名,值为字段值)
                            F.map_from_entries(F.struct(F.map_keys(F.struct(addr)), F.map_values(F.struct(addr)))),
                            lambda kv: kv.value.isNotNull()
                        ),
                        lambda kv: F.concat_ws(": ", kv.key, kv.value)
                    )
                )
        ),
        "; "
    )
)

# 查看结果
processed_df.select("id", "formatted_addresses").show(truncate=False)

执行结果:

+---+----------------------------------------------------------------------------------------------------+
|id |formatted_addresses                                                                                 |
+---+----------------------------------------------------------------------------------------------------+
|1  |street: 123 Main St, city: New York, postalCode: 10001, state: NY, country: USA; city: Los Angeles, postalCode: 90001, country: USA|
|2  |street: 456 Oak Ave, state: CA, country: USA                                                        |
+---+----------------------------------------------------------------------------------------------------+

方式二:自定义UDF(易读性优先)

如果内置函数的逻辑过于复杂,可使用自定义UDF处理单个地址,代码更直观:

from pyspark.sql.functions import udf

# 定义处理单个地址的函数:忽略空值字段,拼接成字符串
def format_single_address(address):
    if not address:
        return ""
    address_parts = []
    for field_name, field_value in address.items():
        if field_value is not None:
            address_parts.append(f"{field_name}: {field_value}")
    return ", ".join(address_parts)

# 注册UDF
format_address_udf = udf(format_single_address, StringType())

# 应用UDF并拼接地址
processed_df_udf = df.withColumn(
    "formatted_addresses",
    F.array_join(F.transform(F.col("addresses"), format_address_udf), "; ")
)

# 查看结果
processed_df_udf.select("id", "formatted_addresses").show(truncate=False)

方案对比

  • 内置函数方式:无需序列化/反序列化,性能更高,适合大数据量场景。
  • UDF方式:代码逻辑更易理解,适合需要复杂自定义规则的场景。

内容的提问来源于stack exchange,提问作者Jatin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:35:20