如何用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
相关产品推荐
相关产品推荐

