A séma from_json következtetése és fejlesztése pipeline-okban

Fontos

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

A Lakeflow-folyamatokban az from_json SQL-függvény automatikusan kikövetkeztetheti és fejlesztheti a JSON-blobok sémáját anélkül, hogy explicit sémát ad meg.

Hogyan működik a(z) from_json a folyamatokban

Az from_json SQL-függvény egy JSON-sztringoszlopot elemez, és egy strukturált értéket ad vissza. Ha egy folyamaton kívül használják, explicit módon meg kell adnia a visszaadott érték sémáját az schema argumentum használatával. Ha egy folyamatban használják, engedélyezheti a sémakövetkeztetést és az evolúciót, amely automatikusan kezeli a visszaadott érték sémáját. Ez a funkció leegyszerűsíti a kezdeti beállítást (különösen akkor, ha a séma ismeretlen) és a folyamatban lévő műveleteket, amikor a séma gyakran változik. Tetszőleges JSON-blobokat dolgoz fel olyan streamelési adatforrásokból, mint az Auto Loader, a Kafka vagy a Kinesis.

Kifejezetten, amikor egy adatfolyamatban használják, az from_json SQL-függvény sémakövetkeztetése és fejlődése az alábbiakat kínálhatja:

  • Új mezők észlelése bejövő JSON-rekordokban (beleértve a beágyazott JSON-objektumokat is)
  • A mezőtípusok következtetése és megfeleltetése a megfelelő Spark-adattípusokhoz
  • A séma automatikus fejlesztése új mezők elhelyezésére
  • Az aktuális sémának nem megfelelő adatok automatikus kezelése

Szintaxis: A séma automatikus következtetése és továbbfejlesztése

Ha egy folyamatban szeretné engedélyezni a sémakövetkeztetést from_json , állítsa a sémát NULL értékre, és adja meg a schemaLocationKey beállítást. Ez lehetővé teszi a séma következtetését és nyomon követését.

SQL

from_json(jsonStr, NULL, map("schemaLocationKey", "<uniqueKey>” [, otherOptions]))

Python

from_json(jsonStr, None, {"schemaLocationKey": "<uniqueKey>”[, otherOptions]})

A lekérdezések több kifejezéssel is from_json rendelkezhetnek, de mindegyik kifejezésnek egyedinek schemaLocationKeykell lennie. A schemaLocationKey folyamatonkénti egyedinek is kell lennie.

SQL

SELECT
  value,
  from_json(value, NULL, map('schemaLocationKey', 'keyX')) parsedX,
  from_json(value, NULL, map('schemaLocationKey', 'keyY')) parsedY
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')

Python

(spark.readStream
    .format("cloudFiles")
    .option("cloudFiles.format", "text")
    .load("/databricks-datasets/nyctaxi/sample/json/")
    .select(
      col("value"),
      from_json(col("value"), None, {"schemaLocationKey": "keyX"}).alias("parsedX"),
      from_json(col("value"), None, {"schemaLocationKey": "keyY"}).alias("parsedY"))
)

Szintaxis: Rögzített séma

Ha ehelyett egy adott sémát szeretne kikényszeríteni, az alábbi from_json szintaxissal elemezheti a JSON-sztringet a sémával:

from_json(jsonStr, schema, [, options])

Ez a szintaxis bármely Azure Databricks-környezetben használható, köztük az adatfolyamatokban is. További információ itt.

Sémakövetkeztetés

from_json A JSON-adatoszlopok első kötegéből következtet a sémára, és belsőleg indexeli azt a schemaLocationKey (kötelező) alapján.

Ha a JSON-sztring egyetlen objektum (például), {"id": 123, "name": "John"} egy STRUCT típusú sémára következtet, from_jsonés hozzáad egy rescuedDataColumn értéket a mezők listájához.

STRUCT<id LONG, name STRING, _rescued_data STRING>

Ha azonban a JSON-sztring felső szintű tömböt (például ["id": 123, "name": "John"]) tartalmaz, akkor from_json a TÖMBöt egy STRUCT-be csomagolja. Ez a módszer lehetővé teszi a következtetett sémával nem kompatibilis adatok mentését. Lehetősége van arra, hogy a tömbértékeket külön sorokba bontsa lefelé.

STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>

Sémakövetkeztetés felülbírálása sématippekkel

Igény szerint megadhatja schemaHints , hogy befolyásolja, hogyan from_json következtethet az oszlop típusára. Ez akkor hasznos, ha tudja, hogy egy oszlop egy adott adattípusból áll, 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 az SQL-séma specifikációjának szintaxisával. A sématippek szemantikája megegyezik az Automatikus betöltő sématippekkel. Például:

SELECT
-- The JSON `{"a": 1}` will treat `a` as a BIGINT
from_json(data, NULL, map('schemaLocationKey', 'w', 'schemaHints', '')),
-- The JSON `{"a": 1}` will treat `a` as a STRING
from_json(data, NULL, map('schemaLocationKey', 'x', 'schemaHints', 'a STRING')),
-- The JSON `{"a": {"b": 1}}` will treat `a` as a MAP<STRING, BIGINT>
from_json(data, NULL, map('schemaLocationKey', 'y', 'schemaHints', 'a MAP<STRING, BIGINT'>)),
-- The JSON `{"a": {"b": 1}}` will treat `a` as a STRING
from_json(data, NULL, map('schemaLocationKey', 'z', 'schemaHints', 'a STRING')),
FROM STREAM READ_FILES(...)

Ha a JSON-sztring egy legfelső szintű TÖMBöt tartalmaz, az egy STRUCT-be van csomagolva. Ezekben az esetekben a rendszer sématippeket alkalmaz a TÖMB sémára a burkolt STRUCT helyett. Vegyük például egy legfelső szintű tömböt tartalmazó JSON-sztringet, például:

[{"id": 123, "name": "John"}]

A kikövetkezett ARRAY-séma egy STRUCT-be van csomagolva:

STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>

Az adattípus id módosításához adja meg a sémamutatót STRING-ként element.id. Új, DOUBLE típusú oszlop hozzáadásához adja meg a element.new_col DOUBLE-t. A fenti tippek miatt a legfelső szintű JSON-tömb sémája a következő lesz:

struct<value array<id STRING, name STRING, new_col DOUBLE>, _rescued_data STRING>

Fejlessze tovább a sémát schemaEvolutionMode

from_json észleli az új oszlopok hozzáadását az adatok feldolgozása során. Új from_json mező észlelésekor a rendszer frissíti a következtetett sémát a legújabb sémával úgy, hogy új oszlopokat egyesít a séma végéhez. A meglévő oszlopok adattípusai változatlanok maradnak. A sémafrissítés után a folyamat automatikusan újraindul a frissített sémával.

from_json A sémafejlődéshez az alábbi módokat támogatja, amelyeket az opcionális schemaEvolutionMode beállítással állít be. Ezek a módok összhangban vannak az automatikus betöltővel.

schemaEvolutionMode Új oszlop olvasásának viselkedése
addNewColumns (alapértelmezett) Az adatfolyam meghiúsult. A rendszer új oszlopokat ad hozzá a sémához. A meglévő oszlopok nem fejlesztik az adattípusokat.
rescue A séma soha nem fejlődik, és a stream nem hiúsul meg sémamódosítások miatt. Minden új oszlop a mentett adatoszlopban lesz rögzítve.
failOnNewColumns Az adatfolyam meghiúsult. A Stream csak akkor indul újra, ha a schemaHints-k frissítésekre kerülnek, vagy eltávolítják a jogsértő adatokat.
none Nem fejleszti a sémát, az új oszlopok figyelmen kívül lesznek hagyva, és az adatok mentése csak akkor történik meg, ha a rescuedDataColumn beállítás be van állítva. A stream nem hiúsul meg sémamódosítások miatt.

Például:

SELECT
-- If a new column appears, the pipeline will automatically add it to the schema:
from_json(a, NULL, map('schemaLocationKey', 'w', 'schemaEvolutionMode', 'addNewColumns')),
-- If a new column appears, the pipeline will add it to the rescued data column:
from_json(b, NULL, map('schemaLocationKey', 'x', 'schemaEvolutionMode', 'rescue')),
-- If a new column appears, the pipeline will ignore it:
from_json(c, NULL, map('schemaLocationKey', 'y', 'schemaEvolutionMode', 'none')),
-- If a new column appears, the pipeline will fail:
from_json(d, NULL, map('schemaLocationKey', 'z', 'schemaEvolutionMode', 'failOnNewColumns')),
FROM STREAM READ_FILES(...)

Mentett adat oszlop

A rendszer automatikusan hozzáad egy mentett adatoszlopot a sémához._rescued_data Az rescuedDataColumn opció beállításával átnevezheti az oszlopot. Például:

from_json(jsonStr, None, {"schemaLocationKey": "keyX", "rescuedDataColumn": "my_rescued_data"})

Ha a mentett adatoszlopot választja, a rendszer a kihagyás helyett azokat az oszlopokat menti, amelyek nem felelnek meg a kikövetkeztetett sémának. Ez egy adattípus eltérés, egy hiányzó oszlop a sémában, vagy az oszlopnév kis- és nagybetű különbsége miatt fordulhat elő.

Sérült rekordok kezelése

A hibásan formázott és nem elemezhető rekordok tárolásához adjon hozzá egy _corrupt_record oszlopot sémamutatók beállításával, például az alábbi példában:

CREATE STREAMING TABLE bronze AS
  SELECT
    from_json(value, NULL,
      map('schemaLocationKey', 'nycTaxi',
          'schemaHints', '_corrupt_record STRING',
          'columnNameOfCorruptRecord', '_corrupt_record')) jsonCol
  FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')

A sérült rekordoszlop átnevezéséhez állítsa be a columnNameOfCorruptRecord beállítást.

A JSON-elemző három módot támogat a sérült rekordok kezelésére:

Üzemmód Description
PERMISSIVE Sérült rekordok esetén a hibásan formázott sztringet egy által konfigurált columnNameOfCorruptRecord mezőbe helyezi, és a hibásan formázott mezőket a következőre nullállítja be: . A sérült rekordok megőrzéséhez beállíthat egy sztring típusú columnNameOfCorruptRecord nevű mezőt egy felhasználó által definiált sémában. Ha egy séma nem rendelkezik a mezővel, a rendszer az elemzés során elveti a sérült rekordokat. Séma következtetésekor az elemző implicit módon hozzáad egy columnNameOfCorruptRecord mezőt a kimeneti sémához.
DROPMALFORMED Figyelmen kívül hagyja a sérült rekordokat.
A DROPMALFORMED mód használata rescuedDataColumn esetén az adattípus eltérései nem okoznak rekordok eldobását. A rendszer csak sérült rekordokat elvet, például hiányos vagy hibásan formázott JSON-rekordokat.
FAILFAST Kivételt eredményez, ha a feldolgozó sérült rekordokkal találkozik.
Ha a FAILFAST módot a rescuedDataColumn használatával használja, az adattípus-eltérések nem eredményeznek hibát. Csak a sérült rekordok okoznak hibákat, például hiányos vagy hibásan formázott JSON-t.

Tekintse meg a from_json kimenet egyik mezőjét

from_json a séma a folyamat végrehajtása során következtetésre jut. Ha egy alsóbb rétegbeli lekérdezés egy from_json mezőre hivatkozik, mielőtt a from_json függvény sikeresen végrehajtotta volna legalább egyszer, a mező nem oldódik fel, és a lekérdezést kihagyja. Az alábbi példában a rendszer kihagyja az ezüsttáblás lekérdezés elemzését, amíg a from_json bronz lekérdezésben lévő függvény nem hajtja végre és nem következteti ki a sémát.

CREATE STREAMING TABLE bronze AS
  SELECT
    from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
  FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')

CREATE STREAMING TABLE silver AS
  SELECT jsonCol.VendorID, jsonCol.total_amount
  FROM bronze

Ha a from_json függvényre és az arra következtetendő mezőkre ugyanabban a lekérdezésben hivatkozik, az elemzés a következő példához hasonlóan meghiúsulhat:

CREATE STREAMING TABLE bronze AS
  SELECT
    from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
  FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
  WHERE jsonCol.total_amount > 100.0

Ezt kijavíthatja úgy, hogy a from_json mezőre mutató hivatkozást egy alsóbb rétegbeli lekérdezésbe helyezi át (például a fenti bronz/ezüst példát).) Másik lehetőségként megadhatja schemaHints , hogy mely mezők tartalmazzák a hivatkozott from_json mezőket. Például:

CREATE STREAMING TABLE bronze AS
  SELECT
    from_json(value, NULL, map('schemaLocationKey', 'nycTaxi', 'schemaHints', 'total_amount DOUBLE')) jsonCol
  FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
  WHERE jsonCol.total_amount > 100.0

Példák: A séma automatikus következtetése és továbbfejlesztése

Ez a szakasz példakódot tartalmaz az automatikus sémakikövetkeztetés és sémaevolúció from_json használatával történő engedélyezéséhez a folyamatokban.

Streaming tábla készítése felhőbeli objektumtárolóból

Az alábbi példa szintaxissal read_files hoz létre streamelési táblázatot a felhőobjektum-tárolóból.

SQL

CREATE STREAMING TABLE bronze AS
  SELECT
    from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
  FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')

Python

@dp.table(comment="from_json autoloader example")
def bronze():
  return (
    spark.readStream
         .format("cloudFiles")
         .option("cloudFiles.format", "text")
         .load("/databricks-datasets/nyctaxi/sample/json/")
         .select(from_json(col("value"), None, {"schemaLocationKey": "nycTaxi"}).alias("jsonCol"))
)

Streamelési tábla létrehozása a Kafkából

Az alábbi példa a read_kafka szintaxist használja egy streamelési tábla létrehozásához a Kafka esetében.

SQL

CREATE STREAMING TABLE bronze AS
  SELECT
    value,
    from_json(value, NULL, map('schemaLocationKey', 'keyX')) jsonCol,
  FROM READ_KAFKA(
    bootstrapSevers => '<server:ip>',
    subscribe => 'events',
    "startingOffsets", "latest"
)

Python

@dp.table(comment="from_json kafka example")
def bronze():
  return (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "<server:ip>")
         .option("subscribe", "<topic>")
         .option("startingOffsets", "latest")
         .load()
         .select(col(“value”), from_json(col(“value”), None, {"schemaLocationKey": "keyX"}).alias("jsonCol"))
)

Példák: Rögzített séma

Például egy rögzített sémát használó from_json kód, lásd from_json a függvényt.

FAQs

Ez a szakasz a függvény sémakövető és evolúciós támogatásával from_json kapcsolatos gyakori kérdésekre ad választ.

Mi a különbség a from_json és a parse_json között?

A parse_json függvény a JSON-sztring VARIANT értékét adja vissza.

A VARIANT rugalmas és hatékony módot kínál a félig strukturált adatok tárolására. Ez megkerüli a sémakövetkeztetést és az evolúciót azáltal, hogy teljesen távolodik a szigorú típusoktól. Ha azonban íráskor szeretne sémát kényszeríteni (például viszonylag szigorú sémával rendelkezik), from_json jobb megoldás lehet.

Az alábbi táblázat a következők közötti from_jsonparse_jsonkülönbségeket ismerteti:

Funkció Használati esetek Availability
from_json Sémafejlődés a séma fenntartásával from_json . Ez akkor hasznos, ha:
  • Szeretné kikényszeríteni az adatsémát (például minden sémamódosítást át kell vizsgálnia annak megőrzése előtt).
  • Optimalizálni szeretné a tárterületet, és alacsony lekérdezési késést és költséget igényel.
  • Nem egyező típusú adatokon szeretne meghiúsulni.
  • Részleges eredményeket szeretne kinyerni sérült JSON-rekordokból, és tárolni a hibásan formázott rekordot az _corrupt_record oszlopban. Ezzel szemben a VARIANT-betöltés érvénytelen JSON-hibát ad vissza.
Csak adatfolyamokban érhető el sémakövetkeztetéssel és -evolúcióval
parse_json A VARIANT különösen alkalmas olyan adatok tartására, amelyeket nem kell sémába illeszteni. Például:
  • Az adatokat részben strukturáltnak szeretné tartani, mert rugalmasak.
  • A séma túl gyorsan változik ahhoz, hogy gyakori streamhibák és újraindítások nélkül illeszkedjen a sémába.
  • A nem egyező típusú adatokon nem szeretne meghiúsulni. (A VARIANT-betöltés mindig sikeres érvényes JSON-rekordok esetén, még akkor is, ha típuseltérések vannak.)
  • A felhasználók nem szeretnék kezelni a mentett adatoszlopot, amely olyan mezőket tartalmaz, amelyek nem felelnek meg a sémának.
A folyamatokon belül és kívül is elérhető

Használhatok from_json sémakövetkezési és evolúciós szintaxist a folyamatokon kívül?

Nem, a from_json sémakövetkeztetési és -evolúciós szintaxist nem használhatod pipeline-okon kívül.

Hogyan férhetek hozzá a from_json által leképezett séma adataihoz?

Tekintse meg a célstreamelési tábla sémáját.

Át tudok adni from_json egy sémát, és evolúciót is végezhetek?

Nem, nem adhat át from_json sémát, és nem is végezhet evolúciót. Azonban megadhat séma tippeket, hogy felülírja a from_json által meghatározott mezők egy részét vagy mindegyikét.

Mi történik a sémával, ha a tábla teljesen frissül?

A rendszer törli a táblához társított sémahelyeket, és a séma újrakövetkeztethető az alapoktól.