You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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解决:

  1. 定义对应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()
  1. 转换为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 09:26:15