Inférer et faire évoluer le schéma en utilisant from_json dans les pipelines
Aperçu
Cette fonctionnalité est en Aperçu public.
Dans les LakeFlow Pipelines, la fonction SQL from_json peut inférer et faire évoluer le schéma des objets blob JSON automatiquement, sans que vous n'ayez à fournir de schéma explicite.
Fonctionnement de from_json dans les pipelines
La fonction SQL from_json analyse une colonne de chaîne JSON et renvoie une valeur de struct. Lorsqu'il est utilisé en dehors d'un pipeline, vous devez explicitement fournir le schéma de la valeur renvoyée à l'aide de l'argument schema. Lorsqu'il est utilisé dans un pipeline, vous pouvez activer l'inférence et l'évolution du schéma, ce qui gère automatiquement le schéma de la valeur renvoyée. Cette fonctionnalité simplifie à la fois la configuration initiale (surtout lorsque le schéma est inconnu) et les Opérations continues lorsque le schéma change fréquemment. Il traite des objets BLOB JSON arbitraires provenant de sources de données de streaming telles qu'Auto Loader, Kafka ou Kinesis.
Plus précisément, lorsqu'elles sont utilisées dans un pipeline, l'inférence et l'évolution de schémas pour la fonction SQL from_json peuvent :
- Détecter les nouveaux champs dans les enregistrements JSON entrants (y compris les objets JSON imbriqués)
- Inférez les types de champs et mappez-les aux types de données Spark appropriés
- Faire évoluer automatiquement le schéma pour prendre en charge les nouveaux champs.
- Gérer automatiquement les données non conformes au schéma actuel
Syntaxe : déduire et faire évoluer automatiquement le schéma
Pour activer l'inférence de schéma avec from_json dans un pipeline, définissez le schéma sur NULL et spécifiez l'option schemaLocationKey. Cela lui permet d'inférer et de suivre le schéma.
- SQL
- Python
from_json(jsonStr, NULL, map("schemaLocationKey", "<uniqueKey>” [, otherOptions]))
from_json(jsonStr, None, {"schemaLocationKey": "<uniqueKey>”[, otherOptions]})
Une query peut avoir plusieurs expressions from_json, mais chaque expression doit avoir un schemaLocationKey unique. Le schemaLocationKey doit également être unique par pipeline.
- SQL
- Python
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')
(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"))
)
Syntaxe : Schéma corrigé
Si vous souhaitez plutôt appliquer un schéma particulier, vous pouvez utiliser la syntaxe from_json suivante pour analyser la chaîne JSON à l'aide de ce schéma :
from_json(jsonStr, schema, [, options])
Cette syntaxe peut être utilisée dans n’importe quel environnement Databricks, y compris les pipelines. Plus d'informations sont disponibles ici.
Inférence du schéma
from_json déduit le schéma du premier batch de colonnes de données JSON et l'indexe en interne par son schemaLocationKey (obligatoire).
Si la chaîne JSON est un objet unique (par exemple, {"id": 123, "name": "John"}), from_json infère un schéma de type STRUCT et ajoute un rescuedDataColumn à la liste des champs.
STRUCT<id LONG, name STRING, _rescued_data STRING>
Cependant, si la chaîne JSON a un tableau de niveau supérieur (tel que ["id": 123, "name": "John"]), alors from_json encapsule le TABLEAU dans un STRUCT. Cette approche permet la récupération des données incompatibles avec le schéma inféré. Vous avez la possibilité de déployer les valeurs du tableau en lignes distinctes en aval.
STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>
Ignorer l'inférence de schéma à l'aide d'indices de schéma
Vous pouvez éventuellement fournir schemaHints pour influencer la manière dont from_json infère le type d'une colonne. Ceci est utile lorsque vous savez qu'une colonne est d'un type de données spécifique, ou si vous souhaitez choisir un type de données plus général (par exemple, un double au lieu d'un entier). Vous pouvez fournir un nombre arbitraire d'indications pour les types de données de colonne à l'aide de la syntaxe de spécification de schéma SQL. La sémantique des indications de schéma est la même que celle des indications de schéma de l'Auto Loader. Par exemple :
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(...)
Lorsqu'une chaîne JSON contient un tableau de premier niveau, elle est encapsulée dans une STRUCT. Dans ces cas, les indications de schéma sont appliquées au schéma ARRAY au lieu de la STRUCT encapsulée. Par exemple, considérez une chaîne JSON avec un tableau de premier niveau, tel que :
[{"id": 123, "name": "John"}]
Le schéma ARRAY inféré est enveloppé dans un STRUCT :
STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>
Pour modifier le type de données de id, spécifiez l'indication de schéma comme element.id STRING. Pour ajouter une nouvelle colonne de type DOUBLE, spécifiez element.new_col DOUBLE. En raison de ces indications, le schéma pour le tableau JSON de niveau supérieur devient :
struct<value array<id STRING, name STRING, new_col DOUBLE>, _rescued_data STRING>
Faire évoluer le schéma en utilisant schemaEvolutionMode
from_json détecte l'ajout de nouvelles colonnes lorsqu'il traite vos données. Lorsque from_json détecte un nouveau champ, il met à jour le schéma inféré avec le dernier schéma en fusionnant les nouvelles colonnes à la fin du schéma. Les types de données des colonnes existantes restent inchangés. Après la mise à jour du schéma, le pipeline redémarre automatiquement avec le schéma mis à jour.
from_json prend en charge les modes suivants pour l'évolution des schémas, que vous définissez à l'aide du paramètre facultatif schemaEvolutionMode. Ces modes sont compatibles avec Auto Loader.
| Comportement lors de la lecture d'une nouvelle colonne |
|---|---|
| Le Stream échoue. De nouvelles colonnes sont ajoutées au schéma. Les types de données des colonnes existantes n'évoluent pas. |
| Le schéma n'évolue jamais et le Stream ne tombe pas en panne en raison des modifications de schéma. Toutes les nouvelles colonnes sont enregistrées dans la colonne de données récupérées. |
| Le Stream échoue. Le Stream ne redémarre pas à moins que les |
| Ne fait pas évoluer le schéma, les nouvelles colonnes sont ignorées et les données ne sont pas récupérées à moins que l'option |
Par exemple :
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(...)
Colonne de données récupérée
Une colonne de données récupérées est automatiquement ajoutée à votre schéma en tant que _rescued_data. Vous pouvez renommer la colonne en définissant l’option rescuedDataColumn. Par exemple :
from_json(jsonStr, None, {"schemaLocationKey": "keyX", "rescuedDataColumn": "my_rescued_data"})
Lorsque vous choisissez d'utiliser la colonne de données récupérées, toutes les colonnes qui ne correspondent pas au schéma inféré sont récupérées au lieu d'être ignorées. Cela peut se produire en raison d'une incompatibilité de type de données, d'une colonne manquante dans le schéma ou d'une différence de casse de nom de colonne.
Gérer les enregistrements corrompus
Pour stocker les enregistrements mal formés et qui ne peuvent pas être analysés, ajoutez une colonne _corrupt_record en définissant des indications de schéma, comme dans l'exemple suivant :
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')
Pour renommer la colonne d'enregistrement corrompu, définissez l'option columnNameOfCorruptRecord.
Le parseur JSON prend en charge trois modes de gestion des enregistrements corrompus :
Mode | Description |
|---|---|
| Pour les enregistrements corrompus, place la chaîne mal formée dans un champ configuré par |
| Ignore les enregistrements corrompus. Lorsque vous utilisez le mode |
| Lève une exception lorsque l'analyseur rencontre des enregistrements corrompus. Lorsque vous utilisez le mode |
Référencer un champ dans la sortie from_json
from_json infère le schéma lors de l'exécution du pipeline. Si une requête en aval fait référence à un champ from_json avant que la fonction from_json ne se soit exécutée avec succès au moins une fois, le champ n'est pas résolu et la requête est ignorée. Dans l'exemple suivant, l'analyse de la requête de la table argent sera ignorée jusqu'à ce que la fonction from_json dans la requête bronze se soit exécutée et ait inféré le schéma.
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
Si la fonction from_json et les champs qu'elle infère sont mentionnés dans la même query, l'analyse peut échouer comme dans l'exemple suivant :
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
Vous pouvez corriger cela en déplaçant la référence au champ from_json dans une query en aval (comme dans l'exemple bronze/silver ci-dessus). Alternativement, vous pouvez spécifier schemaHints qui contiennent les champs from_json référencés. Par exemple :
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
Exemples : Déduire et faire évoluer automatiquement le schéma
Cette section fournit un exemple de code pour activer l'inférence et l'évolution automatiques du schéma à l'aide de from_json dans les pipelines.
Créez une table de streaming à partir du stockage d'objets cloud
L'exemple suivant utilise la syntaxe read_files pour créer une table de streaming à partir du stockage d'objets cloud.
- SQL
- Python
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')
@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"))
)
Créer une table de streaming à partir de Kafka
L'exemple suivant utilise la syntaxe read_kafka pour créer une table de streaming à partir de Kafka.
- SQL
- Python
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"
)
@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"))
)
Exemples : Schéma fixe
Pour un exemple de code utilisant from_json avec un schéma fixe, consultez la fonctionfrom_json.
FAQ
Cette section répond aux questions fréquemment posées sur la prise en charge de l'inférence et de l'évolution du schéma dans la fonction from_json.
Quelle est la différence entre from_json et parse_json?
La fonction parse_json renvoie une valeur VARIANT de la chaîne JSON.
VARIANT offre un moyen flexible et efficace de stocker des données semi-structurées. Cela contourne l'inférence et l'évolution du schéma en éliminant complètement les types stricts. Cependant, si vous souhaitez appliquer un schéma au moment de l'écriture (par exemple, parce que vous avez un schéma relativement strict), from_json pourrait être une meilleure option.
Le tableau suivant décrit les différences entre from_json et parse_json:
Fonction | Cas d’utilisation | Disponibilité |
|---|---|---|
| L'évolution des schémas avec
| Disponible avec l'inférence et l'évolution de schéma uniquement dans les pipelines |
| VARIANT est particulièrement bien adapté pour contenir des données qui n'ont pas besoin d'être schématisées. Par exemple :
| Disponible dans et en dehors des pipelines |
Puis-je utiliser la syntaxe from_json d'inférence et d'évolution de schéma en dehors des pipelines ?
Non, vous ne pouvez pas utiliser la syntaxe d'inférence et d'évolution de schéma from_json en dehors des pipelines.
Comment puis-je accéder au schéma déduit par from_json?
Affichez le schéma de la table de streaming cible.
Puis-je passer from_json un schéma et faire également de l'évolution ?
Non, vous ne pouvez pas transmettre from_json un schéma et également effectuer une évolution. Cependant, vous pouvez fournir des indications de schéma pour remplacer certains ou tous les champs déduits par from_json.
Que se passe-t-il au schéma si la table est entièrement actualisée ?
Les emplacements de schéma associés à la table sont effacés et le schéma est réinféré à partir de zéro.