如何将PySpark DataFrame的JSON字符串列转为JSON对象避免Cosmos DB转义
PySpark 写入Cosmos DB JSON字符串转义问题最优解决方案
出现转义的根本原因是DataFrame中存储的是字符串类型值,Spark写入Cosmos DB时会将整段字符串作为单个字符串值序列化,因此会对内部双引号自动添加转义符。将字符串列解析为Spark原生复杂类型后写入即可解决该问题。
具体实现步骤
- 第一步:定义与JSON字符串结构匹配的Schema
你提供的JSON字符串为数组结构,内部包含一个JSON对象,对应Schema定义如下:
from pyspark.sql.types import StructType, StructField, StringType, ArrayType inner_schema = StructType([ StructField("Date", StringType(), nullable=True), StructField("Ids", StringType(), nullable=True), StructField("col2", StringType(), nullable=True), StructField("col3", StringType(), nullable=True), StructField("value3", StringType(), nullable=True), StructField("currency", StringType(), nullable=True), StructField("status", StringType(), nullable=True), StructField("tag", StringType(), nullable=True), StructField("Info", StringType(), nullable=True) ]) json_schema = ArrayType(inner_schema, nullable=True)
如果你的JSON结构存在嵌套或者其他数据类型,只需要对应调整Schema定义即可。
- 第二步:使用
from_json函数解析JSON字符串列
将原有的JSON字符串列替换为解析后的复杂类型列:
from pyspark.sql.functions import from_json # 假设存储JSON字符串的列名为json_col,替换为你实际的列名 df_parsed = df.withColumn("json_col", from_json("json_col", json_schema))
- 第三步:直接写入Cosmos DB即可
解析后的json_col为Spark原生Array复杂类型,写入Cosmos DB时会被序列化为原生JSON结构,不会再出现双引号被转义的问题。
特殊场景适配
如果你的JSON字符串结构不固定,无法提前定义Schema,可以使用schema_of_json函数动态推断Schema,不需要手动编写:
from pyspark.sql.functions import schema_of_json, lit # 取一条JSON样本动态生成Schema sample_json = '[{"Date":"11/13/2020 3:23:21 PM","Ids":"[]","col2":"abc","col3":"","value3":"[]","currency":"","status":"Active","tag":"[]","Info":"[]"}]' dynamic_schema = schema_of_json(lit(sample_json)) df_parsed = df.withColumn("json_col", from_json("json_col", dynamic_schema))
内容的提问来源于stack exchange,提问作者Prabhu Venkatesan
相关产品推荐
相关产品推荐

