XML-fájlok olvasása és írása

Fontos

Ez a funkció a nyilvános előzetes verzióban érhető el.

A bővíthető korrektúranyelv (XML) az adatok szöveges formátumban történő formázására, tárolására és megosztására szolgáló korrektúranyelv. Az adatok szerializálására vonatkozó szabályokat határoz meg a dokumentumoktól az tetszőleges adatstruktúrákig.

Azure Databricks támogatja az XML-t az Apache Sparkkal való olvasáshoz és íráshoz, beleértve az automatikus sémakövetkeztetést és -fejlesztést, a sorcímkék konfigurációját, az XSD-érvényesítést és az SQL-kifejezéseket, példáulfrom_xml. A natív XML-támogatás működik az Auto Loaderrel, a read_files és a COPY INTO használatával, külső JAR-fájlok használata nélkül.

Prerequisites

Az XML-fájlformátumok támogatásához a Databricks Runtime 14.3-s vagy újabb verziója szükséges.

Beállítások

Az XML-adatforrások konfigurálásához használja a(z) .option() és .options()DataFrameReader és DataFrameWriter metódusait. A támogatott beállítások teljes listájáért tekintse meg DataFrameReader az XML-beállításokat és DataFrameWriter az XML-beállításokat.

XML-rekordok elemzése

Az XML-specifikáció egy jól formázott struktúrát ad meg. Ez a specifikáció azonban nem képez le azonnal táblázatos formátumot. Meg kell adnia a rowTag beállítást annak jelzésére, hogy melyik XML-elem felel meg egy DataFrameRow-nek. Az rowTag elem lesz a legfelső szintű struct. A(z) rowTag gyerekelemei a legfelső szintű struct mezőivé válnak.

Megadhatja a rekord sémáját, vagy hagyhatja, hogy az automatikusan következtethető legyen. Mivel az elemző csak az elemeket vizsgálja, a rendszer kiszűri a rowTag DTD-t és a külső entitásokat.

Az alábbi példák egy XML-fájl sémakövetkeztetését és elemzését szemléltetik különböző rowTag-beállítások használatával:

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)

Olvassa be az XML-fájlt a(z) rowTag beállítással, mint "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'
)

Kimenet:

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}] |
+----------------------------------------------------------------------------------------+

Olvassa be az XML-fájlt rowTag használatával mint "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'
)

Kimenet:

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     |
+----+------+--------------------------------+------+

XML-rekordok ellenőrzése XSD-vel

Igény szerint az egyes sorszintű XML-rekordokat egy XML-sémadefinícióval (XSD) ellenőrizheti. Az XSD-fájl meg van adva a rowValidationXSDPath beállításban. Az XSD egyébként nem befolyásolja a megadott vagy kikövetkezett sémát. Az ellenőrzés meghiúsuló rekordja "sérültként" van megjelölve, és a beállítás szakaszban leírt sérült rekordkezelési mód alapján lesz kezelve.

A XSDToSchema használatával kinyerhet egy Spark DataFrame-sémát egy XSD-fájlból. Csak egyszerű, összetett és szekvenciatípusokat támogat, és csak az alapszintű XSD-funkciókat támogatja.

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))

Az alábbi táblázat az XSD-adattípusok Spark-adattípusokra való konvertálását mutatja be:

XSD-adattípusok Spark-adattípusok
boolean BooleanType
decimal DecimalType
unsignedLong DecimalType(38, 0)
double DoubleType
float FloatType
byte ByteType
short, unsignedByte ShortType
integer, negativeInteger, nonNegativeInteger, nonPositiveIntegerpositiveIntegerunsignedShort IntegerType
long, unsignedInt LongType
date DateType
dateTime TimestampType
Others StringType

Beágyazott XML elemzése

Egy meglévő DataFrame-adatkeret sztring típusú oszlopában lévő XML-adatok az schema_of_xml és from_xml segítségével elemezhetők, amelyek a sémát és az elemzési eredményeket új struct oszlopokként adják vissza. Az argumentumként schema_of_xmlfrom_xml átadott XML-adatoknak egyetlen jól formázott XML-rekordnak kell lenniük.

XML séma

A(z) schema_of_xml használatával XML-karakterláncból következtethet a Spark-sémára. Adja át az eredményt az from_xml XML-oszlopok elemzéséhez.

Szintaxis: schema_of_xml(xmlStr [, options])

Argument Required Description
xmlStr Yes Sztringkifejezés, amely egyetlen jól formázott XML-rekordot határoz meg.
options No A MAP<STRING,STRING> direktívákat meghatározó literál.

Egy olyan karakterláncot ad vissza, amely egy struktúra definícióját tartalmazza n sztringmezőkkel, ahol az oszlopnevek az XML-elemből és az attribútumnevekből származnak. A mezőértékek a származtatott formázott SQL-típusokat tartják.

XML-ből

XML-rekordokat tartalmazó KARAKTERLÁNC-oszlop strukturált szerkezetbe való elemzésére használható from_xml . Adjon meg egy sémát közvetlenül, vagy használja a kimenetet schema_of_xml.

Szintaxis: from_xml(xmlStr, schema [, options])

Argument Required Description
xmlStr Yes Sztringkifejezés, amely egyetlen jól formázott XML-rekordot határoz meg.
schema Yes Egy STRING kifejezés vagy a schema_of_xml függvény meghívása.
options No Egy MAP<STRING,STRING> direktívákat megadó literál.

A sémadefiníciónak megfelelő mezőneveket és típusokat tartalmazó szerkezetet ad vissza. A sémát vesszővel tagolt oszlopnévként és adattípus-párként kell definiálni, például CREATE TABLE. A Beállítások szakaszban látható legtöbb lehetőség a következő kivételekkel alkalmazható:

  • rowTag: Mivel csak egy XML-rekord van, a rowTag beállítás nem alkalmazható.
  • mode (alapértelmezett: PERMISSIVE): Lehetővé teszi a sérült rekordok elemzés közbeni kezelését.
    • PERMISSIVE: Amikor sérült rekorddal találkozik, a hibás karakterláncot a columnNameOfCorruptRecord által konfigurált mezőbe helyezi, és a hibás mezőket null értékre állítja. A sérült rekordok megőrzéséhez beállíthat egy columnNameOfCorruptRecord nevű sztring típusú mezőt egy felhasználó által definiált sémában. Ha egy séma nem rendelkezik a mezővel, az elemzés során a sérült rekordokat elveti. Séma következtetése esetén implicit módon hozzáad egy columnNameOfCorruptRecord mezőt egy kimeneti sémához.
    • FAILFAST: Kivételt eredményez, ha sérült rekordoknak felel meg.

Példák

XML-sztringet tartalmazó oszlop elemzéséhez használja a séma kikövetkeztetéséhez a(z) schema_of_xml elemet, majd adja át a(z) from_xml elemnek:

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()

Beágyazott XML elemzése az SQL-ben:

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>')
);

Konvertálás XML- és DataFrame-struktúrák között

A DataFrame és az XML közötti szerkezeti különbségek miatt bizonyos átalakítási szabályok vonatkoznak az XML-adatok DataFrame formátumba, illetve a DataFrame formátumból XML-adatokká történő átalakítására. Vegye figyelembe, hogy az attribútumok kezelése letiltható a beállítással excludeAttribute.

Konvertálás XML-ről DataFrame-re

XML olvasása során Azure Databricks az XML-elemeket és -attribútumokat a DataFrame-mezőkre képezi le az alábbi szabályok szerint.

Az attribútumok a címsorelőtaggal attributePrefixrendelkező mezőkké alakulnak.

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

Ez a következő sémát hozza létre:

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

Az attribútum(ok) vagy gyermekelem(ek)et tartalmazó elem karakteradatait a rendszer a valueTag mezőbe elemzi. Ha a karakteradatok több előfordulása is előfordul, a valueTag mező típussá array lesz konvertálva.

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

Ez a következő sémát hozza létre:

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)

Konvertálás DataFrame-ről XML-re

DataFrame XML-be írásakor bizonyos beágyazott struktúrák speciális kezelést igényelnek a DataFrame és az XML-adatmodellek közötti különbségek miatt.

Ha egy DataFrame olyan mezőt tartalmaz ArrayType , amelynek elemtípusa is ArrayType, az XML-be írásakor további beágyazási szint jön létre, amely nem jelenik meg az XML-fájlok ciklikus bemásolásakor. Ez csak az XML-en kívüli forrású DataFrame-eket érinti – az XML-fájlok olvasása és írása megőrzi az eredeti struktúrát.

Egy DataFrame például a következő sémával:

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

és a következő adatok:

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

a következő XML-kimenetet hozza létre:

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

A DataFrame névtelen tömbje elemének nevét a arrayElementName beállítás adja meg (alapértelmezett: item).

A mentett adatoszlop engedélyezése

A mentett adatoszlop biztosítja, hogy az ETL során soha ne veszítsen el adatokat. Rögzíti azokat az adatokat, amelyeket nem elemeztek, mert egy rekord egy vagy több mezője az alábbi problémák egyikével rendelkezik:

  • Hiányzik a megadott sémából.
  • Nem egyezik a megadott séma adattípusával.
  • A megadott séma mezőneveivel kis- és nagybetű érzékenység tekintetében nem egyezik meg.

A mentett adatoszlop JSON-dokumentumként lesz visszaadva, amely a mentett oszlopokat és a rekord forrásfájl-elérési útját tartalmazza.

A helyreállított adatok oszlopának engedélyezéséhez beolvasáskor állítsa a(z) rescuedDataColumn beállítást egy oszlopnévre:

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'
)

Ha el szeretné távolítani a forrásfájl elérési útját a mentett adatoszlopból, állítsa be a következőt:

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

Az XML-elemző három módot támogat rekordok elemzésekor: PERMISSIVE, DROPMALFORMEDés FAILFAST. Ha a(z) rescuedDataColumn opcióval együtt használják, az adattípus-eltérések nem okozzák a rekordok eldobását DROPMALFORMED módban, és nem dobnak hibát FAILFAST módban. Csak a sérült rekordok (hiányos vagy helytelenül formázott XML) kerülnek eldobásra, vagy okoznak hibát.

Séma következtetése és fejlesztése az Automatikus betöltővel

A témakör részletes ismertetését és az alkalmazható lehetőségeket Sémakövetkeztetés és -fejlesztés konfigurálása az Automatikus betöltőcímű témakörben találja. Az Automatikus betöltőt úgy konfigurálhatja, hogy automatikusan észlelje a betöltött XML-adatok sémáját, így anélkül inicializálhatja a táblákat, hogy explicit módon deklarálja az adatsémát, és új oszlopok bevezetésekor továbbfejleszti a táblázatsémát. Ez szükségtelenné teszi a sémamódosítások manuális nyomon követését és alkalmazását az idő függvényében.

Az automatikus betöltő sémakövetkeztetése alapértelmezés szerint a típuseltérések miatti sémafejlődési problémák elkerülésére törekszik. Olyan formátumok esetén, amelyek nem kódolnak adattípusokat (JSON, CSV és XML), az Automatikus betöltő minden oszlopot sztringként von le, beleértve az XML-fájlok beágyazott mezőit is. Az Apache Spark DataFrameReader eltérő viselkedést használ a sémakövetkeztetéshez, és mintaadatok alapján választja ki az XML-források oszlopainak adattípusait. Ha engedélyezni szeretné ezt a viselkedést az Auto Loaderrel, állítsa be a cloudFiles.inferColumnTypes opciót true-re.

Az Automatikus betöltő észleli az új oszlopok hozzáadását az adatok feldolgozása során. Amikor az Auto Loader új oszlopot észlel, a stream egy UnknownFieldExceptionhibakóddal leáll. Mielőtt a stream ezt a hibát észleli, az Automatikus betöltő sémakövetkeztetést hajt végre a legújabb mikro-adatkötegen, és frissíti a séma helyét a legújabb sémával az új oszlopoknak a séma végéhez való egyesítésével. A meglévő oszlopok adattípusai változatlanok maradnak. Az Auto Loader különböző módokat támogat a sémafejlődéshez, amelyet az opciók között cloudFiles.schemaEvolutionModeállíthat be.

A sémautalások segítségével érvényesítheti azokat a sémaadatokat, amelyeket ismer és elvár egy következtetett sémában. Ha tudja, hogy egy oszlop adott adattípusú, vagy ha általánosabb adattípust (például egész szám helyett dupla) szeretne választani, tetszőleges számú tippet adhat meg az oszlop adattípusaihoz sztringként az SQL-séma specifikációjának szintaxisával. Ha a mentett adatoszlop engedélyezve van, az olyan mezők, amelyek neve eltér a séma elnevezési stílusától, a _rescued_data oszlopba kerülnek. Ezt a viselkedést úgy módosíthatja, hogy a(z) readerCaseSensitive beállítást false értékre állítja; ebben az esetben az Auto Loader a kis- és nagybetűket figyelmen kívül hagyva olvassa be az adatokat.

Usage

Az alábbi példák a Wanderbricks-adatkészlet használatával szemléltetik az XML-fájlok olvasását és írását a Spark DataFrame API és az SQL használatával.

XML olvasása és írása

A DataFrame API-val a Wanderbricks-véleményeket XML-be írhatja, és visszaolvassa őket.

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")

Adatok olvasásakor manuálisan is megadhatja a sémát:

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")

XML olvasása és írása SQL-lel

Az SQL DDL használatával hozzon létre egy táblát egy XML-fájlból. Azure Databricks automatikusan oszloptípusokat következtet.

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

A DDL-ben oszlopneveket és -típusokat is megadhat. Ebben az esetben a séma nem lesz automatikusan kikövetkeztetve.

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");

XML betöltése COPY INTO

Xml-fájlok COPY INTO betöltése felhőbeli tárolóból Delta-táblába.

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 olvasása sorérvényesítéssel

Ezzel a rowValidationXSDPath beállítással olvasás közben ellenőrizheti az egyes sorokat egy XSD-sémán.

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 betöltése automatikus betöltővel

Az Automatikus betöltő használatával folyamatosan betölthet XML-fájlokat a felhőbeli tárolóból egy Delta-táblába automatikus sémakövetkezéssel és -fejlesztéssel.

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")

További erőforrások