Well, it's not trivial as it would seems. First, your approach is not meant for Spark, unless you're working with very little data (and so, you don't need Spark) and you're better off using pure Python like you tried. Using collect() fetch all data on the driver which would not work with large data.
The distributed approach for this is as follows:
- infer schema on part of your JSON data (unless you want to do it manually - which is tedious)
- transform your dataframe with this schema to have access to named attributes
- select attributes as needed and back to JSON
I tried to decompose as much as I could here:
from pyspark.sql.types import IntegerType, StructType, StringType, StructField
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
import json
spark = SparkSession.builder.getOrCreate()
sc = spark.sparkContext
# Create input data
data = [json.dumps({"A": "2022-01-26T14:21:32.214+0000", "B": 69, "C": {"CA": 42, "CB": "Hello"}, "D": "XD"})]
df = spark.createDataFrame(data, "string").toDF("colA")
df.show()
+-----------------------------------------------------------------------------------------+
|colA |
+-----------------------------------------------------------------------------------------+
|{"A": "2022-01-26T14:21:32.214+0000", "B": 69, "C": {"CA": 42, "CB": "Hello"}, "D": "XD"}|
+-----------------------------------------------------------------------------------------+
# Infer schema - infering on first 10 rows
s = df.select(F.col("colA").alias("s")).rdd.map(lambda x: x.s).take(10)
schema = spark.read.json(sc.parallelize(s)).schema
print(schema)
# StructType(List(StructField(A,StringType,true),StructField(B,LongType,true),StructField(C,StructType(List(StructField(CA,LongType,true),StructField(CB,StringType,true))),true),StructField(D,StringType,true)))
# read JSON string with schema
new_df = df.withColumn("colB", F.from_json("colA", schema))
new_df.show(truncate=False)
+-----------------------------------------------------------------------------------------+---------------------------------------------------+
|colA |colB |
+-----------------------------------------------------------------------------------------+---------------------------------------------------+
|{"A": "2022-01-26T14:21:32.214+0000", "B": 69, "C": {"CA": 42, "CB": "Hello"}, "D": "XD"}|{2022-01-26T14:21:32.214+0000, 69, {42, Hello}, XD}|
+-----------------------------------------------------------------------------------------+---------------------------------------------------+
# Finally ...
new_df.select(F.to_json(F.struct("colB.A", "colB.B", "colB.D")).alias("colC")).show(truncate=False)
+----------------------------------------------------+
|colC |
+----------------------------------------------------+
|{"A":"2022-01-26T14:21:32.214+0000","B":69,"D":"XD"}|
+----------------------------------------------------+