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

PySpark中如何将JSON字符串转为JSON对象?含多JSON结构列处理

Got it, let's solve this problem step by step. You've got a DataFrame where the test column holds string-formatted JSONs of four different (but partially overlapping) structures, and you need to convert those into structured JSON objects (Spark structs, which act like native JSON objects in Spark). Here's how to do it:

Step 1: Define Schemas for Each JSON Type

First, we need to create schemas that match each of the four JSON structures. You can infer these schemas directly from sample JSON strings of each type:

# Sample JSON strings for each of your 4 types (replace with your actual sample data)
sample_purchase = '{"PurchaseActivity": {"transaction_id": "txn_123", "amount": 49.99, "timestamp": "2023-11-05T14:30:00Z"}}'
sample_interaction = '{"Interaction": {"user_id": "user_456", "action": "view", "page": "product_detail"}}'
sample_other_type_1 = '{"OtherType1": {"event_id": "evt_789", "status": "completed", "metadata": {"source": "app"}}}'
sample_other_type_2 = '{"OtherType2": {"order_id": "ord_012", "items": ["item_a", "item_b"], "shipping_address": {"city": "New York"}}}'

# Infer schemas for each type using Spark's JSON reader
purchase_schema = spark.read.json(sc.parallelize([sample_purchase])).schema
interaction_schema = spark.read.json(sc.parallelize([sample_interaction])).schema
other1_schema = spark.read.json(sc.parallelize([sample_other_type_1])).schema
other2_schema = spark.read.json(sc.parallelize([sample_other_type_2])).schema
Step 2: Classify Rows by JSON Type

Next, we'll add a column to identify which JSON structure each row uses. Using a regex is the most reliable way to capture the top-level key, even if there's leading whitespace:

from pyspark.sql.functions import col, when, from_json, regexp_extract

# Extract the top-level JSON key to classify each row
df = df.withColumn(
    "json_type",
    regexp_extract(col("test"), r'^\s*{"([^"]+)":', 1)
)
Step 3: Parse Each Row with the Correct Schema

Now we can use Spark's from_json function to convert the string into a structured JSON object, applying the matching schema based on the json_type we just created:

# Parse the string JSON into a structured Spark struct
df_parsed = df.withColumn(
    "parsed_json",
    when(col("json_type") == "PurchaseActivity", from_json(col("test"), purchase_schema))
    .when(col("json_type") == "Interaction", from_json(col("test"), interaction_schema))
    .when(col("json_type") == "OtherType1", from_json(col("test"), other1_schema))
    .when(col("json_type") == "OtherType2", from_json(col("test"), other2_schema))
    .otherwise(None)  # Handle unrecognized types as needed
)
Step 4: (Optional) Flatten or Access Nested Fields

If you want to work with individual fields instead of the full struct, you can extract them using dot notation. For example, to pull common fields across all structures:

from pyspark.sql.functions import coalesce

df_flattened = df_parsed.select(
    "test",
    "json_type",
    # Coalesce common fields from all structures (adjust based on your actual shared fields)
    coalesce(
        col("parsed_json.PurchaseActivity.timestamp"),
        col("parsed_json.Interaction.timestamp"),
        col("parsed_json.OtherType1.timestamp"),
        col("parsed_json.OtherType2.timestamp")
    ).alias("common_timestamp"),
    # Extract unique fields for each type
    col("parsed_json.PurchaseActivity.amount").alias("purchase_amount"),
    col("parsed_json.Interaction.action").alias("interaction_action")
)
Verification

You can check the result by printing the schema and sample rows:

df_parsed.printSchema()
df_parsed.show(10, truncate=False)

内容的提问来源于stack exchange,提问作者Arvind-MSFT

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:40:18