Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
A séma
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:
|
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:
|
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.