在Databricks中从S3读取复杂JSON并转为字符串的方法
问题:Databricks读取S3多层JSON并转为字符串
我在Databricks环境中工作,需要从S3读取多层结构的JSON文件并转换为字符串。示例JSON结构如下:
{ "id": "123", "details":[ { "name": "Bob", "address": "123 street" }, { "name": "Amy", "address": "XYZ street" } ], "docType": "File", "collections": ["a","b","c"] }
我通过以下代码挂载S3并读取JSON为Spark DataFrame:
aws_s3_bucket = 'my_bucket' mount_name = '/mnt/test' source_url = 's3a://%s' %(aws_s3_bucket) dbutils.fs.mount(source_url,mount_name) file_path = "/dummyKey/dummyFile.json" df = spark.read.option("multiline","true").json(mount_name + file_path).cache()
该DataFrame包含id、details、docType、collections四个列。尝试使用toJson()函数及Pandas的to_json方法将其转为JSON字符串时,均出现如下错误:
the queries from raw JSON/CSV files are disallowed when the referenced columns only include the internal corrupt record column
想寻求可行方案:如何从S3读取该JSON文件并转换为字符串?是否可以直接读取为字符串而不转为DataFrame?
解决方案
方案一:直接读取为字符串(无需转DataFrame)
可以绕过Spark DataFrame,直接读取原始文件内容为字符串,适合不需要对JSON结构做处理的场景:
# 方法1:使用dbutils.fs读取 file_content = dbutils.fs.head(mount_name + file_path) # 方法2:通过/dbfs前缀访问挂载路径,用Python原生open读取 with open('/dbfs' + mount_name + file_path, 'r') as f: file_content = f.read()
说明:
dbutils.fs.head()适合读取常规大小的JSON文件,大文件建议分块读取- Databricks中挂载的路径可以通过
/dbfs前缀映射为本地文件系统路径,支持Python原生文件操作
方案二:修复DataFrame转JSON的错误
如果必须通过Spark DataFrame处理后再转字符串,错误根源是Spark解析JSON时可能隐式生成了_corrupt_record列,可通过显式指定Schema解决:
- 定义对应JSON结构的Schema:
from pyspark.sql.types import StructType, StructField, StringType, ArrayType schema = StructType([ StructField("id", StringType(), True), StructField("details", ArrayType(StructType([ StructField("name", StringType(), True), StructField("address", StringType(), True) ])), True), StructField("docType", StringType(), True), StructField("collections", ArrayType(StringType()), True) ]) # 带Schema读取JSON,避免自动解析错误 df = spark.read.option("multiline","true").schema(schema).json(mount_name + file_path).cache()
- 转换为JSON字符串:
# 方式1:用Spark原生toJSON() json_rdd = df.toJSON() # 原文件是单条记录,取第一个元素即可 full_json_str = json_rdd.collect()[0] # 方式2:转Pandas DataFrame后处理 pandas_df = df.toPandas() # orient='records'会生成数组格式,去掉首尾[]还原单条记录 full_json_str = pandas_df.to_json(orient='records')[1:-1]
说明:
- 显式指定Schema能确保Spark只解析预期字段,避免生成
_corrupt_record列 - 单条JSON记录场景下,两种转换方式都需要做简单处理,还原原始JSON结构
内容的提问来源于stack exchange,提问作者Joe
相关产品推荐
相关产品推荐

