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

Spark解析JSON:获取子字段值为true的父字段名称

Spark 解决方案:提取带指定标记的字段名

需求说明

处理给定的嵌套JSON数据,提取unique_id: true对应的父字段名,以及ignore: true对应的父字段名,最终输出指定格式的结果表格。

实现步骤与代码

以下是基于PySpark的实现代码:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, collect_list, explode

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

# 示例JSON数据(实际场景可替换为文件路径,如spark.read.json("/path/to/data.json"))
json_data = """
{    
    "properties": {
        "student_id": {
            "type": "string",
            "unique_id": true
        },
        "status": {
            "type": "boolean"
        },
        "name": {
            "type": "string"
        },
        "phone_number": {
            "type": "string",
            "ignore": true
        },
        "e_mail": {
            "type": "string",
            "ignore": true
        },
        "address": {
            "type": "string"
        }
    },
    "subjects": [
        "science",
        "english"
    ]
}
"""

# 加载JSON数据为DataFrame
df = spark.read.json(spark.sparkContext.parallelize([json_data]))

# 拆分properties字段的键值对,得到字段名和字段详情
properties_df = df.select(explode(col("properties")).alias("field_name", "field_details"))

# 收集所有unique_id为true的字段名
unique_id_fields = properties_df.filter(col("field_details.unique_id") == True)\
    .agg(collect_list("field_name").alias("unique_id"))

# 收集所有ignore为true的字段名
ignore_fields = properties_df.filter(col("field_details.ignore") == True)\
    .agg(collect_list("field_name").alias("ignore"))

# 合并两个结果集
result_df = unique_id_fields.crossJoin(ignore_fields)

# 按指定格式输出结果
# 先查看Spark原生输出
result_df.show(truncate=False)

# 生成示例要求的表格格式
pandas_result = result_df.toPandas()
print("\n指定格式输出:")
print("unique_id     | ignore                 |")
print("----------------------------------------")
unique_str = f"[{pandas_result['unique_id'][0][0]}]" if len(pandas_result['unique_id'][0]) > 0 else "[]"
ignore_str = f"[{', '.join(pandas_result['ignore'][0])}]" if len(pandas_result['ignore'][0]) > 0 else "[]"
print(f"{unique_str:<13} | {ignore_str:<23} |")

代码说明

  1. 数据加载:将JSON数据加载为Spark DataFrame,支持本地/分布式文件路径输入。
  2. 字段拆分:使用explode函数将properties这个Map类型字段拆分为字段名和字段详情两列,实现逐个字段的遍历检查。
  3. 筛选聚合:分别筛选出unique_id和ignore标记为true的字段,通过collect_list将符合条件的字段名聚合为列表。
  4. 结果合并:通过crossJoin将两个聚合结果合并为一行,保证结果结构符合要求。
  5. 格式输出:先通过Spark的show方法查看原生结果,再转为Pandas DataFrame后按照指定的表格格式打印输出,兼容空结果场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 13:07:37