JSON 文件

Spark SQL 可以自动推断 JSON 数据集的模式(Schema)并将其加载为 DataFrame。这种转换可以通过在 JSON 文件上使用 SparkSession.read.json 来完成。

请注意,作为 JSON 文件 提供的文件并非典型的 JSON 文件。每一行必须包含一个独立的、自包含的有效 JSON 对象。有关更多信息,请参阅 JSON Lines 文本格式,也称为换行符分隔的 JSON

对于常规的多行 JSON 文件,请将 multiLine 参数设置为 True

# spark is from the previous example.
sc = spark.sparkContext

# A JSON dataset is pointed to by path.
# The path can be either a single text file or a directory storing text files
path = "examples/src/main/resources/people.json"
peopleDF = spark.read.json(path)

# The inferred schema can be visualized using the printSchema() method
peopleDF.printSchema()
# root
#  |-- age: long (nullable = true)
#  |-- name: string (nullable = true)

# Creates a temporary view using the DataFrame
peopleDF.createOrReplaceTempView("people")

# SQL statements can be run by using the sql methods provided by spark
teenagerNamesDF = spark.sql("SELECT name FROM people WHERE age BETWEEN 13 AND 19")
teenagerNamesDF.show()
# +------+
# |  name|
# +------+
# |Justin|
# +------+

# Alternatively, a DataFrame can be created for a JSON dataset represented by
# an RDD[String] storing one JSON object per string
jsonStrings = ['{"name":"Yin","address":{"city":"Columbus","state":"Ohio"}}']
otherPeopleRDD = sc.parallelize(jsonStrings)
otherPeople = spark.read.json(otherPeopleRDD)
otherPeople.show()
# +---------------+----+
# |        address|name|
# +---------------+----+
# |[Columbus,Ohio]| Yin|
# +---------------+----+
完整的示例代码可在 Spark 仓库的“examples/src/main/python/sql/datasource.py”中找到。

Spark SQL 可以自动推断 JSON 数据集的模式并将其加载为 Dataset[Row]。这种转换可以通过在 Dataset[String] 或 JSON 文件上使用 SparkSession.read.json() 来完成。

请注意,作为 JSON 文件 提供的文件并非典型的 JSON 文件。每一行必须包含一个独立的、自包含的有效 JSON 对象。有关更多信息,请参阅 JSON Lines 文本格式,也称为换行符分隔的 JSON

对于常规的多行 JSON 文件,请将 multiLine 选项设置为 true

// Primitive types (Int, String, etc) and Product types (case classes) encoders are
// supported by importing this when creating a Dataset.
import spark.implicits._

// A JSON dataset is pointed to by path.
// The path can be either a single text file or a directory storing text files
val path = "examples/src/main/resources/people.json"
val peopleDF = spark.read.json(path)

// The inferred schema can be visualized using the printSchema() method
peopleDF.printSchema()
// root
//  |-- age: long (nullable = true)
//  |-- name: string (nullable = true)

// Creates a temporary view using the DataFrame
peopleDF.createOrReplaceTempView("people")

// SQL statements can be run by using the sql methods provided by spark
val teenagerNamesDF = spark.sql("SELECT name FROM people WHERE age BETWEEN 13 AND 19")
teenagerNamesDF.show()
// +------+
// |  name|
// +------+
// |Justin|
// +------+

// Alternatively, a DataFrame can be created for a JSON dataset represented by
// a Dataset[String] storing one JSON object per string
val otherPeopleDataset = spark.createDataset(
  """{"name":"Yin","address":{"city":"Columbus","state":"Ohio"}}""" :: Nil)
val otherPeople = spark.read.json(otherPeopleDataset)
otherPeople.show()
// +---------------+----+
// |        address|name|
// +---------------+----+
// |[Columbus,Ohio]| Yin|
// +---------------+----+
完整的示例代码可在 Spark 仓库的“examples/src/main/scala/org/apache/spark/examples/sql/SQLDataSourceExample.scala”中找到。

Spark SQL 可以自动推断 JSON 数据集的模式并将其加载为 Dataset<Row>。这种转换可以通过在 Dataset<String> 或 JSON 文件上使用 SparkSession.read().json() 来完成。

请注意,作为 JSON 文件 提供的文件并非典型的 JSON 文件。每一行必须包含一个独立的、自包含的有效 JSON 对象。有关更多信息,请参阅 JSON Lines 文本格式,也称为换行符分隔的 JSON

对于常规的多行 JSON 文件,请将 multiLine 选项设置为 true

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;

// A JSON dataset is pointed to by path.
// The path can be either a single text file or a directory storing text files
Dataset<Row> people = spark.read().json("examples/src/main/resources/people.json");

// The inferred schema can be visualized using the printSchema() method
people.printSchema();
// root
//  |-- age: long (nullable = true)
//  |-- name: string (nullable = true)

// Creates a temporary view using the DataFrame
people.createOrReplaceTempView("people");

// SQL statements can be run by using the sql methods provided by spark
Dataset<Row> namesDF = spark.sql("SELECT name FROM people WHERE age BETWEEN 13 AND 19");
namesDF.show();
// +------+
// |  name|
// +------+
// |Justin|
// +------+

// Alternatively, a DataFrame can be created for a JSON dataset represented by
// a Dataset<String> storing one JSON object per string.
List<String> jsonData = Arrays.asList(
        "{\"name\":\"Yin\",\"address\":{\"city\":\"Columbus\",\"state\":\"Ohio\"}}");
Dataset<String> anotherPeopleDataset = spark.createDataset(jsonData, Encoders.STRING());
Dataset<Row> anotherPeople = spark.read().json(anotherPeopleDataset);
anotherPeople.show();
// +---------------+----+
// |        address|name|
// +---------------+----+
// |[Columbus,Ohio]| Yin|
// +---------------+----+
完整的示例代码可在 Spark 仓库的“examples/src/main/java/org/apache/spark/examples/sql/JavaSQLDataSourceExample.java”中找到。

Spark SQL 可以自动推断 JSON 数据集的模式并使用 read.json() 函数将其加载为 DataFrame,该函数从 JSON 文件目录中加载数据,其中文件的每一行都是一个 JSON 对象。

请注意,作为 JSON 文件 提供的文件并非典型的 JSON 文件。每一行必须包含一个独立的、自包含的有效 JSON 对象。有关更多信息,请参阅 JSON Lines 文本格式,也称为换行符分隔的 JSON

对于常规的多行 JSON 文件,请将命名参数 multiLine 设置为 TRUE

# A JSON dataset is pointed to by path.
# The path can be either a single text file or a directory storing text files.
path <- "examples/src/main/resources/people.json"
# Create a DataFrame from the file(s) pointed to by path
people <- read.json(path)

# The inferred schema can be visualized using the printSchema() method.
printSchema(people)
## root
##  |-- age: long (nullable = true)
##  |-- name: string (nullable = true)

# Register this DataFrame as a table.
createOrReplaceTempView(people, "people")

# SQL statements can be run by using the sql methods.
teenagers <- sql("SELECT name FROM people WHERE age >= 13 AND age <= 19")
head(teenagers)
##     name
## 1 Justin
完整的示例代码可在 Spark 仓库的“examples/src/main/r/RSparkSQLExample.R”中找到。
CREATE TEMPORARY VIEW jsonTable
USING org.apache.spark.sql.json
OPTIONS (
  path "examples/src/main/resources/people.json"
)

SELECT * FROM jsonTable

数据源选项

JSON 的数据源选项可以通过以下方式设置

属性名称默认值含义范围
timeZone (spark.sql.session.timeZone 配置的值) 设置用于格式化 JSON 数据源或分区值中时间戳的时区 ID 字符串。支持以下 timeZone 格式:
  • 基于区域的时区 ID:格式应为 'area/city',例如 'America/Los_Angeles'。
  • 时区偏移量(Zone offset):应采用 '(+|-)HH:mm' 格式,例如 '-08:00' 或 '+01:00'。此外,'UTC' 和 'Z' 也作为 '+00:00' 的别名受到支持。
不建议使用 'CST' 等其他短名称,因为它们可能存在歧义。
读取/写入
primitivesAsString false 将所有原始值推断为字符串类型。 读取
prefersDecimal false 将所有浮点值推断为十进制类型(decimal type)。如果值无法容纳在 decimal 类型中,则将其推断为双精度浮点数(doubles)。 读取
allowComments false 忽略 JSON 记录中 Java/C++ 风格的注释。 读取
allowUnquotedFieldNames false 允许使用未加引号的 JSON 字段名称。 读取
allowSingleQuotes true 除双引号外,还允许使用单引号。 读取
allowNumericLeadingZeros false 允许数字的前导零(例如 00012)。 读取
allowBackslashEscapingAnyCharacter false 允许使用反斜杠转义机制对所有字符进行引用。 读取
mode PERMISSIVE 允许在解析过程中处理损坏记录的模式。
  • PERMISSIVE:当遇到损坏的记录时,将畸形的字符串放入由 columnNameOfCorruptRecord 配置的字段中,并将畸形字段设置为 null。为了保留损坏的记录,用户可以在自定义模式中设置一个名为 columnNameOfCorruptRecord 的字符串类型字段。如果模式中没有该字段,则在解析过程中会丢弃损坏的记录。在推断模式时,它会在输出模式中隐式添加一个 columnNameOfCorruptRecord 字段。
  • DROPMALFORMED:忽略整个损坏的记录。JSON 内置函数不支持此模式。
  • FAILFAST:遇到损坏记录时抛出异常。
读取
columnNameOfCorruptRecord (spark.sql.columnNameOfCorruptRecord 配置的值) 允许重命名由 PERMISSIVE 模式创建的包含畸形字符串的新字段。这会覆盖 spark.sql.columnNameOfCorruptRecord。 读取
dateFormat yyyy-MM-dd 设置指示日期格式的字符串。自定义日期格式遵循 日期时间模式 中的格式。此项适用于日期类型。 读取/写入
timestampFormat yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] 设置指示时间戳格式的字符串。自定义日期格式遵循 日期时间模式 中的格式。此项适用于时间戳类型。 读取/写入
timestampNTZFormat yyyy-MM-dd'T'HH:mm:ss[.SSS] 设置指示不带时区的时间戳格式的字符串。自定义日期格式遵循 日期时间模式。这适用于不带时区的时间戳类型;注意,在读写此数据类型时,不支持时区偏移和时区组件。 读取/写入
enableDateTimeParsingFallback 如果时间解析策略具有旧版设置或未提供自定义日期或时间戳模式,则启用此项。 允许在值与设置的模式不匹配时,回退到解析日期和时间戳的向后兼容(Spark 1.x 和 2.0)行为。 读取
multiLine false 每个文件解析一条记录,该记录可能跨越多行。JSON 内置函数会忽略此选项。 读取
allowUnquotedControlChars false 是否允许 JSON 字符串包含未加引号的控制字符(ASCII 值小于 32 的字符,包括制表符和换行符)。 读取
encoding multiLine 设置为 true(读取时)或 UTF-8(写入时)时自动检测。 对于读取,允许强制设置 JSON 文件的标准基本或扩展编码之一。例如 UTF-16BE、UTF-32LE。对于写入,指定保存的 JSON 文件的编码(字符集)。JSON 内置函数会忽略此选项。 读取/写入
lineSep \r\r\n\n(用于读取),\n(用于写入) 定义用于解析的行分隔符。JSON 内置函数会忽略此选项。 读取/写入
samplingRatio 1.0 定义用于模式推断的输入 JSON 对象的比例。 读取
dropFieldIfAllNull false 在模式推断期间,是否忽略全是 null 值或空数组的列。 读取
locale en-US 以 IETF BCP 47 格式的语言标签设置区域设置(Locale)。例如,locale 在解析日期和时间戳时使用。 读取
allowNonNumericNumbers true 允许 JSON 解析器将一组“非数值”(NaN)标记识别为合法的浮点数值。
  • +INF:用于正无穷大,以及 +InfinityInfinity 的别名。
  • -INF:用于负无穷大,-Infinity 的别名。
  • NaN:用于其他非数值,例如除以零的结果。
读取
compression (无) 保存到文件时使用的压缩编解码器。这可以是已知的对大小写不敏感的缩写名称之一(none、bzip2、gzip、lz4、snappy 和 deflate)。JSON 内置函数会忽略此选项。 写入
ignoreNullFields (spark.sql.jsonGenerator.ignoreNullFields 配置的值) 生成 JSON 对象时是否忽略 null 字段。 写入
useUnsafeRow (spark.sql.json.useUnsafeRow 配置的值) 在 JSON 解析器中是否使用 UnsafeRow 来表示结构体结果。 读取

其他通用选项可以在 通用文件源选项 中找到。