讀取和寫入 XML 檔案

重要

這項功能處於公開預覽狀態。

可擴充標記語言 (XML)是一種用於格式化、儲存及分享文字格式資料的標記語言。 它會定義一組規則,以串行化從檔到任意數據結構的數據。

Azure Databricks 支援透過 Apache Spark 讀取及寫入 XML,包括自動結構描述推斷與演化、列標籤設定、XSD 驗證,以及如 from_xml 這類 SQL 運算式。 原生 XML 支援可與 Auto Loader、read_files 和 COPY INTO 搭配使用,無需外部 jar 檔。

先決條件

XML 檔案格式支援需要 Databricks Runtime 14.3 及以上版本。

選項

使用 .option() 和 .options() 的 DataFrameReader 與 DataFrameWriter 方法來設定 XML 資料來源。 欲了解完整支援選項清單,請參閱 DataFrameReader XML 選項 與 DataFrameWriter XML 選項。

剖析 XML 記錄

XML 規格規定格式正確的結構。 不過,此規格不會立即對應至表格格式。 您必須指定 rowTag 選項,用來指明對應到 DataFrameRow 的 XML 元素。 元素 rowTag 會變成最上層 struct。 的子項目 rowTag 會成為最上層 struct的欄位。

您可以指定此記錄的架構,或讓它自動推斷。 因為剖析器只會檢查 rowTag 元素,因此會篩選掉 DTD 和外部實體。

下列範例說明如何使用不同的 rowTag 選項來推斷 XML 檔案的結構描述並加以剖析:

Python

xmlString = """
  <reviews>
    <review id="r001">
      <author>Alice</author>
      <rating>5</rating>
      <comment>Amazing stay, highly recommend!</comment>
    </review>
    <review id="r002">
      <author>Bob</author>
      <rating>4</rating>
      <comment>Great location, very comfortable</comment>
    </review>
  </reviews>"""

xmlPath = "/Volumes/<catalog>/<schema>/<volume>/reviews.xml"
dbutils.fs.put(xmlPath, xmlString, True)

程式語言 Scala

val xmlString = """
  <reviews>
    <review id="r001">
      <author>Alice</author>
      <rating>5</rating>
      <comment>Amazing stay, highly recommend!</comment>
    </review>
    <review id="r002">
      <author>Bob</author>
      <rating>4</rating>
      <comment>Great location, very comfortable</comment>
    </review>
  </reviews>"""
val xmlPath = "/Volumes/<catalog>/<schema>/<volume>/reviews.xml"
dbutils.fs.put(xmlPath, xmlString)

帶有 rowTag 選項的 XML 檔案可讀為 "reviews":

Python

df = spark.read.option("rowTag", "reviews").format("xml").load(xmlPath)
df.printSchema()
df.show(truncate=False)

程式語言 Scala

val df = spark.read.option("rowTag", "reviews").xml(xmlPath)
df.printSchema()
df.show(truncate=false)

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews.xml',
  format => 'xml',
  rowTag => 'reviews'
)

輸出:

root
|-- review: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- _id: string (nullable = true)
| | |-- author: string (nullable = true)
| | |-- comment: string (nullable = true)
| | |-- rating: string (nullable = true)

+----------------------------------------------------------------------------------------+
|review                                                                                  |
+----------------------------------------------------------------------------------------+
|[{r001, Alice, Amazing stay, highly recommend!, 5}, {r002, Bob, Great location..., 4}] |
+----------------------------------------------------------------------------------------+

將含有 rowTag 的 XML 檔案讀取為 "review":

Python

df = spark.read.option("rowTag", "review").format("xml").load(xmlPath)
# Infers four top-level fields and parses `review` in separate rows:

程式語言 Scala

val df = spark.read.option("rowTag", "review").xml(xmlPath)
// Infers four top-level fields and parses `review` in separate rows:

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews.xml',
  format => 'xml',
  rowTag => 'review'
)

輸出:

root
|-- _id: string (nullable = true)
|-- author: string (nullable = true)
|-- comment: string (nullable = true)
|-- rating: string (nullable = true)

+----+------+--------------------------------+------+
|_id |author|comment                         |rating|
+----+------+--------------------------------+------+
|r001|Alice |Amazing stay, highly recommend! |5     |
|r002|Bob   |Great location, very comfortable|4     |
+----+------+--------------------------------+------+

使用 XSD 驗證 XML 紀錄

您可以選擇是否使用 XML 架構定義 (XSD) 驗證每個資料列層級的 XML 記錄。 XSD 檔案是在 rowValidationXSDPath 選項中指定的。 XSD 不會以其他方式影響已提供或推斷的架構。 驗證失敗的紀錄會被標記為「損壞」,並依據選項區段所述的損壞記錄處理模式選項進行處理。

您可以使用 XSDToSchema 從 XSD 檔案擷取 Spark DataFrame 架構。 它只支持簡單、複雜和循序類型,而且只支援基本的 XSD 功能。

import org.apache.spark.sql.execution.datasources.xml.XSDToSchema
import org.apache.hadoop.fs.Path

val xsdPath = "/Volumes/<catalog>/<schema>/<volume>/reviews.xsd"
val xsdString = """<?xml version="1.0" encoding="UTF-8" ?>
  <xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema">
    <xs:element name="review">
      <xs:complexType>
        <xs:sequence>
          <xs:element name="author" type="xs:string" />
          <xs:element name="rating" type="xs:integer" />
          <xs:element name="comment" type="xs:string" />
        </xs:sequence>
        <xs:attribute name="id" type="xs:string" use="required" />
      </xs:complexType>
    </xs:element>
  </xs:schema>"""

dbutils.fs.put(xsdPath, xsdString, true)

val schema1 = XSDToSchema.read(xsdString)
val schema2 = XSDToSchema.read(new Path(xsdPath))

下表顯示 XSD 資料類型轉換成 Spark 資料類型:

XSD 數據類型 Spark 資料類型
boolean BooleanType
decimal DecimalType
unsignedLong DecimalType(38, 0)
double DoubleType
float FloatType
byte ByteType
short、unsignedByte ShortType
integer、、negativeIntegernonNegativeInteger、nonPositiveInteger、、positiveInteger、unsignedShort IntegerType
long、unsignedInt LongType
date DateType
dateTime TimestampType
Others StringType

解析巢狀 XML

現有 DataFrame 中字串值資料行內的 XML 資料可以使用 schema_of_xml 和 from_xml 進行剖析,並以新的 struct 資料行傳回結構描述和剖析結果。 作為引數傳遞給 schema_of_xml 和 from_xml 的 XML 資料,必須是單一且格式正確的 XML 記錄。

XML 結構模式

用 schema_of_xml 來從 XML 字串推斷 Spark 架構。 將結果傳遞給 from_xml 以剖析 XML 欄位。

語法:schema_of_xml(xmlStr [, options])

論點 Required Description
xmlStr Yes 一個 STRING 表達式,指定一個格式良好的 XML 記錄。
options No 一個 MAP<STRING,STRING> 字面上的具體指令。

回傳一個字串,包含一個結構體的定義,該結構體包含 n 個欄位,欄位名稱皆源自 XML 元素與屬性名稱。 域值會保存衍生的格式化 SQL 類型。

從 XML 匯入

用於 from_xml 解析包含 XML 紀錄的 STRING 欄位,並建立結構體。 直接提供結構描述,或使用 schema_of_xml 的輸出。

語法:from_xml(xmlStr, schema [, options])

論點 Required Description
xmlStr Yes 一個 STRING 表達式,指定一個格式良好的 XML 記錄。
schema Yes STRING 表達式或呼叫 schema_of_xml 函式。
options No 一個 MAP<STRING,STRING> 字面上的具體指令。

回傳一個帶有欄位名稱與型別與結構定義相符的結構體。 架構必須定義為逗號分隔的資料列名稱和資料類型群組,例如 CREATE TABLE 選項區塊中顯示的大多數選項皆適用,但有以下例外:

  • rowTag:因為只有一個 XML 記錄,因此 rowTag 選項不適用。
  • mode(預設值:PERMISSIVE):允許使用一種在剖析期間處理損毀記錄的模式。
    • PERMISSIVE:當遇到損毀記錄時,會將格式錯誤的字串放入由 columnNameOfCorruptRecord 設定的欄位中,並將格式錯誤的欄位設為 null。 若要保留損毀的記錄,您可以在使用者定義的架構中設定名為 columnNameOfCorruptRecord 的字串類型字段。 如果結構沒有該欄位,則在解析過程中會忽略損壞的記錄。 推斷架構時,它會隱含地在輸出架構中新增 columnNameOfCorruptRecord 字段。
    • FAILFAST:遇到損毀的紀錄時會引發例外。

範例

要解析 XML 字串欄位,請使用 schema_of_xml 來推斷結構,然後傳給 from_xml:

Python

from pyspark.sql.functions import from_xml, schema_of_xml, lit, col

xml_data = """
  <review id="r001">
    <author>Alice</author>
    <rating>5</rating>
    <comment>Amazing stay, highly recommend!</comment>
  </review>
"""

df = spark.createDataFrame([(1, xml_data)], ["review_id", "payload"])
schema = schema_of_xml(df.select("payload").limit(1).collect()[0][0])
parsed = df.withColumn("parsed", from_xml(col("payload"), schema))
parsed.printSchema()
parsed.show()

程式語言 Scala

import org.apache.spark.sql.functions.{from_xml, schema_of_xml, lit}

val xmlData = """
  <review id="r001">
    <author>Alice</author>
    <rating>5</rating>
    <comment>Amazing stay, highly recommend!</comment>
  </review>""".stripMargin

val df = Seq((1, xmlData)).toDF("review_id", "payload")
val schema = schema_of_xml(xmlData)
val parsed = df.withColumn("parsed", from_xml($"payload", schema))
parsed.printSchema()
parsed.show()

若要剖析 SQL 中的內嵌 XML:

SELECT from_xml('
  <review id="r001">
    <author>Alice</author>
    <rating>5</rating>
    <comment>Amazing stay, highly recommend!</comment>
  </review>',
  schema_of_xml('
  <review id="r001">
    <author>Alice</author>
    <rating>5</rating>
    <comment>Amazing stay, highly recommend!</comment>
  </review>')
);

XML 與 DataFrame 結構間的轉換

由於 DataFrame 與 XML 之間的結構差異,因此有一些從 XML 資料轉換為 DataFrame 以及從 DataFrame 轉換為 XML 資料的規則。 請注意,使用 選項 excludeAttribute可以停用處理屬性。

從 XML 轉換成 DataFrame

讀取 XML 時,Azure Databricks 會依照以下規則將 XML 元素與屬性映射到 DataFrame 欄位。

屬性會轉換為帶有標題前綴 attributePrefix的欄位。

<one myOneAttrib="AAAA">
  <two>two</two>
  <three>three</three>
</one>

這會產生以下模式:

root
|-- _myOneAttrib: string (nullable = true)
|-- two: string (nullable = true)
|-- three: string (nullable = true)

包含屬性或子元素的元素中的角色資料會被解析到 valueTag 欄位中。 如果字元資料出現多次,valueTag 欄位會轉換為 array 類型。

<one>
  <two myTwoAttrib="BBBBB">two</two>
  some value between elements
  <three>three</three>
  some other value between elements
</one>

這會產生以下模式:

root
 |-- _VALUE: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- two: struct (nullable = true)
 |    |-- _VALUE: string (nullable = true)
 |    |-- _myTwoAttrib: string (nullable = true)
 |-- three: string (nullable = true)

從 DataFrame 轉換成 XML

當將資料框架寫成 XML 時,由於資料框架與 XML 資料模型的差異,某些巢狀結構需要特殊處理。

如果一個資料框架包含 ArrayType 一個欄位,其元素類型也是 ArrayType,寫入 XML 會產生一個額外的巢狀層級,這是在循環 XML 檔案時不存在的。 這只影響來自 XML 以外的資料框架——讀寫 XML 檔案會保留原始結構。

例如,具有以下結構的資料框架:

|-- a: array (nullable = true)
| |-- element: array (containsNull = true)
| | |-- element: string (containsNull = true)

以及以下資料:

+------------------------------------+
| a|
+------------------------------------+
|[WrappedArray(aa), WrappedArray(bb)]|
+------------------------------------+

產生以下 XML 輸出:

<a>
  <item>aa</item>
</a>
<a>
  <item>bb</item>
</a>

中 DataFrame 未命名陣列的元素名稱是由 選項 arrayElementName 指定 (預設值: item)。

啟用已救出的資料欄位

已救援資料欄位確保你在 ETL 期間不會遺失資料。 它會擷取任何未被解析的資料,因為記錄中的一個或多個欄位存在以下其中一種問題:

  • 從提供的架構中缺席。
  • 不符合所提供結構描述的資料類型。
  • 具有與所提供結構描述中欄位名稱不符的情況。

已修復的資料行會以 JSON 文件的形式傳回,其中包含已修復的資料行,以及記錄的來源檔案路徑。

要啟用已救援的資料欄位,讀取時請將選項設 rescuedDataColumn 為欄位名稱:

Python

df = spark.read.option("rescuedDataColumn", "_rescued_data").format("xml").load("/Volumes/<catalog>/<schema>/<volume>/reviews_xml")

程式語言 Scala

val df = spark.read.option("rescuedDataColumn", "_rescued_data").format("xml").load("/Volumes/<catalog>/<schema>/<volume>/reviews_xml")

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_xml',
  format => 'xml',
  rowTag => 'review',
  rescuedDataColumn => '_rescued_data'
)

若要從已救援的資料欄位移除來源檔案路徑,請設定:

spark.conf.set("spark.databricks.sql.rescuedDataColumn.filePath.enabled", "false")

剖析記錄時,XML 剖析器支援三種模式: PERMISSIVE、 DROPMALFORMED和 FAILFAST。 與 rescuedDataColumn 一起使用時,資料類型不符不會導致記錄在 DROPMALFORMED 模式中遭到捨棄,也不會在 FAILFAST 模式中引發錯誤。 只有損壞的記錄(不完整或格式錯誤的 XML)才會被捨棄或拋出錯誤。

用自動載入器推斷並演化結構

如需本主題和適用選項的詳細討論,請參閱 設定自動載入器中的架構推斷和演進。 您可以設定 Auto Loader 自動偵測已載入 XML 資料的結構描述,讓您在無須明確宣告資料結構描述的情況下初始化資料表,並在新增資料行時演進資料表結構描述。 這樣就不需要在一段時間內手動追蹤和套用架構變更。

依預設,Auto Loader 架構推斷會主動避免因類型不匹配而導致的架構演化問題。 對於未編碼數據類型的格式(JSON、CSV 和 XML),自動載入器會將所有數據行推斷為字串,包括 XML 檔案中的巢狀字段。 Apache Spark DataFrameReader 會針對架構推斷使用不同的行為,根據範例數據選取 XML 來源中數據行的數據類型。 若要使用自動載入器開啟此行為,請將 選項 cloudFiles.inferColumnTypes 設定為 true。

自動加載器在處理您的資料時會偵測到新增新資料行的情況。 當自動載入器偵測到新的數據行時,數據流會以 UnknownFieldException停止。 在資料流出現此錯誤之前,自動載入器會先對最新的微批次資料執行架構推斷,並以最新的架構更新架構位置,將新欄位合併至架構結尾。 現有數據行的數據類型保持不變。 自動載入器支援架構演進的不同模式,您可以在 選項 cloudFiles.schemaEvolutionMode中設定。

您可以使用 結構描述提示,對推斷出的結構描述強制套用您已知且預期的結構描述資訊。 當您知道某個資料行屬於特定的資料類型,或想要選擇較通用的資料類型(例如使用 double 而非 integer)時,您可以使用 SQL 結構描述規格語法,以字串形式提供任意數量的資料行資料類型提示。 啟用救援資料欄時,名稱大小寫與結構描述不一致的欄位會載入到 _rescued_data 欄中。 您可以將 選項 readerCaseSensitive 設定為 false來變更此行為,在此情況下,自動載入器會以不區分大小寫的方式讀取數據。

Usage

以下範例使用 Wanderbricks 資料集 示範使用 Spark DataFrame API 與 SQL 讀寫 XML 檔案。

讀取和寫入 XML

使用 DataFrame API 將 Wanderbricks 評論寫成 XML 格式並回讀。

Python

# Write Wanderbricks reviews to XML
df = spark.read.table("samples.wanderbricks.reviews")
df.write \
  .format("xml") \
  .option("rootTag", "reviews") \
  .option("rowTag", "review") \
  .save("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")

# Read the XML file back
df_read = spark.read \
  .format("xml") \
  .option("rowTag", "review") \
  .load("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")
df_read.show()

程式語言 Scala

// Write Wanderbricks reviews to XML
val df = spark.read.table("samples.wanderbricks.reviews")
df.write
  .format("xml")
  .option("rootTag", "reviews")
  .option("rowTag", "review")
  .save("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")

// Read the XML file back
val dfRead = spark.read
  .format("xml")
  .option("rowTag", "review")
  .xml("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")
dfRead.show()

R

df <- loadDF("/Volumes/<catalog>/<schema>/<volume>/reviews.xml", source = "xml", rowTag = "review")
saveDF(df, "/Volumes/<catalog>/<schema>/<volume>/newreviews.xml", "xml", "overwrite")

讀取資料時,您可以手動指定架構:

Python

from pyspark.sql.types import StructType, StructField, StringType, IntegerType

custom_schema = StructType([
    StructField("_id", StringType(), True),
    StructField("author", StringType(), True),
    StructField("rating", IntegerType(), True),
    StructField("comment", StringType(), True)
])
df = spark.read.options(rowTag='review').xml('/Volumes/<catalog>/<schema>/<volume>/reviews.xml', schema=custom_schema)
df.show()

程式語言 Scala

import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType}

val customSchema = StructType(Array(
  StructField("_id", StringType, nullable = true),
  StructField("author", StringType, nullable = true),
  StructField("rating", IntegerType, nullable = true),
  StructField("comment", StringType, nullable = true)))
val df = spark.read.option("rowTag", "review").schema(customSchema).xml("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")
df.show()

R

customSchema <- structType(
  structField("_id", "string"),
  structField("author", "string"),
  structField("rating", "integer"),
  structField("comment", "string"))

df <- loadDF("/Volumes/<catalog>/<schema>/<volume>/reviews.xml", source = "xml", schema = customSchema, rowTag = "review")
saveDF(df, "/Volumes/<catalog>/<schema>/<volume>/newreviews.xml", "xml", "overwrite")

用 SQL 讀寫 XML

使用 SQL DDL 從 XML 檔案建立資料表。 Azure Databricks 會自動推斷欄位類型。

DROP TABLE IF EXISTS reviews;
CREATE TABLE reviews
USING XML
OPTIONS (path "/Volumes/<catalog>/<schema>/<volume>/reviews.xml", rowTag "review");
SELECT * FROM reviews;

您也可以在 DDL 中指定資料行名稱和類型。 在此情況下,不會自動推斷架構。

DROP TABLE IF EXISTS reviews;

CREATE TABLE reviews (_id string, author string, rating integer, comment string)
USING XML
OPTIONS (path "/Volumes/<catalog>/<schema>/<volume>/reviews.xml", rowTag "review");

使用 COPY INTO 載入 XML

用 COPY INTO 來從雲端儲存載入 XML 檔案到 Delta 表格。

DROP TABLE IF EXISTS reviews;
CREATE TABLE IF NOT EXISTS reviews;

COPY INTO reviews
FROM "/Volumes/<catalog>/<schema>/<volume>/reviews.xml"
FILEFORMAT = XML
FORMAT_OPTIONS ('mergeSchema' = 'true', 'rowTag' = 'review')
COPY_OPTIONS ('mergeSchema' = 'true');

讀取具有數據列驗證的 XML

使用 rowValidationXSDPath 選項,在讀取時依據 XSD 結構描述驗證每一列。

Python

df = (spark.read
    .format("xml")
    .option("rowTag", "review")
    .option("rowValidationXSDPath", xsdPath)
    .load("/Volumes/<catalog>/<schema>/<volume>/reviews.xml"))
df.printSchema()

程式語言 Scala

val df = spark.read
  .option("rowTag", "review")
  .option("rowValidationXSDPath", xsdPath)
  .xml("/Volumes/<catalog>/<schema>/<volume>/reviews.xml")
df.printSchema

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews.xml',
  format => 'xml',
  rowTag => 'review',
  rowValidationXSDPath => '/Volumes/<catalog>/<schema>/<volume>/reviews.xsd'
)

使用自動載入器載入 XML

使用 AutoLoader 持續從雲端儲存將 XML 檔案匯入 Delta 表格,並自動推論結構與演化。

Python

query = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "xml")
  .option("rowTag", "review")
  .option("cloudFiles.inferColumnTypes", True)
  .option("cloudFiles.schemaLocation", schemaPath)
  .option("cloudFiles.schemaEvolutionMode", "rescue")
  .load(inputPath)
  .writeStream
  .option("mergeSchema", "true")
  .option("checkpointLocation", checkPointPath)
  .trigger(availableNow=True)
  .toTable("reviews")
)

程式語言 Scala

val query = spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "xml")
  .option("rowTag", "review")
  .option("cloudFiles.inferColumnTypes", true)
  .option("cloudFiles.schemaLocation", schemaPath)
  .option("cloudFiles.schemaEvolutionMode", "rescue")
  .load(inputPath)
  .writeStream
  .option("mergeSchema", "true")
  .option("checkpointLocation", checkPointPath)
  .trigger(Trigger.AvailableNow())
  .toTable("reviews")

其他資源