PySpark中如何高效解析字符串格式的含JSON列表字段并生成独立DataFrame
Got it, converting to Pandas for JSON parsing is definitely a no-go for large datasets—let's fix this with pure PySpark operations that scale properly. Here's a step-by-step solution using PySpark's native JSON parsing functions to extract your desired DataFrame without sacrificing performance.
Step 1: Define the JSON Schema
First, we need to define the schema that matches the structure inside your model_payload JSON list. Note that some numeric fields are stored as strings in the JSON, so we'll handle type casting later:
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, BooleanType, DoubleType, ArrayType # Define the schema for a single JSON object in the list single_payload_schema = StructType([ StructField("dataset_A", StringType(), nullable=True), StructField("dataset_B", StringType(), nullable=True), StructField("ks_statistic", StringType(), nullable=True), # Stored as string in JSON StructField("pvalue", StringType(), nullable=True), # Stored as string in JSON StructField("rejected_hypothesis", BooleanType(), nullable=True), StructField("target_ks_statistic", StringType(), nullable=True), # Stored as string StructField("target_pvalue", StringType(), nullable=True) # Stored as string ]) # Since model_payload is a list of these objects, wrap in ArrayType payload_schema = ArrayType(single_payload_schema)
Step 2: Parse and Explode the JSON List
Next, we'll use from_json to convert the string column into a structured array, then explode to turn each element in the array into a separate row:
# Start with your original DataFrame (named original_df here) parsed_df = original_df.withColumn( "parsed_payload", F.from_json(F.col("model_payload"), payload_schema) ).withColumn( "exploded_payload", F.explode(F.col("parsed_payload")) # Split array elements into individual rows )
Step 3: Extract and Cast Fields to Match Target Schema
Now we can pull out each field from the exploded JSON structure, cast the string-based numeric fields to DoubleType, and map action_name to the action column as required:
final_df = parsed_df.select( F.col("exploded_payload.dataset_A"), F.col("exploded_payload.dataset_B"), F.col("exploded_payload.ks_statistic").cast(DoubleType()).alias("ks_statistic"), F.col("exploded_payload.pvalue").cast(DoubleType()).alias("pvalue"), F.col("exploded_payload.rejected_hypothesis"), F.col("exploded_payload.target_ks_statistic").cast(DoubleType()).alias("target_ks_statistic"), F.col("exploded_payload.target_pvalue").cast(DoubleType()).alias("target_pvalue"), F.col("action_name").alias("action") )
Key Notes for Robustness
- Handling Malformed JSON: By default,
from_jsonusesPERMISSIVEmode, which returnsnullfor invalid JSON entries. If you want to drop malformed rows entirely, usefrom_json(..., options={"mode": "DROPMALFORMED"}). If you prefer to fail fast on invalid JSON, use{"mode": "FAILFAST"}. - Performance: This approach runs entirely on Spark's distributed engine, so it scales seamlessly with large datasets—no need to collect data to the driver (which is what happens when you convert to Pandas).
内容的提问来源于stack exchange,提问作者m_bash

