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:
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
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) )
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 )
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") )
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

