Ingérer des données en tant que type variant semi-structuré
Dans Databricks Runtime 15.3 et versions ultérieures, vous pouvez utiliser le type VARIANT pour ingérer des données semi-structurées. Cet article décrit le comportement et fournit des exemples de modèles pour ingérer des données depuis un stockage d'objets cloud en utilisant Auto Loader et COPY INTO, des enregistrements streaming depuis Kafka, et des commandes SQL pour créer de nouvelles tables avec des données variantes ou insérer de nouveaux enregistrements en utilisant le type variant. Le tableau suivant récapitule les formats de fichier pris en charge et le support de version de Databricks Runtime :
Format de fichier | Versions de Databricks Runtime prises en charge |
|---|---|
JSON | 15.3 et versions ultérieures |
XML | 16.4 et supérieur |
CSV | 16.4 et supérieur |
Voir Données variantes de query.
Créer une table avec une colonne variante
VARIANT est un type SQL standard dans Databricks Runtime 15.3 et versions ultérieures et est pris en charge par les tables basées sur Delta Lake. Les tables gérées sur Databricks utilisent Delta Lake par default, vous pouvez donc créer une table vide avec une seule colonne VARIANT à l'aide de la syntaxe suivante.
CREATE TABLE table_name (variant_column VARIANT)
Vous pouvez également utiliser une instruction CTAS pour créer une table avec une colonne variante. Utilisez la fonction PARSE_JSON pour analyser les chaînes JSON ou la fonction FROM_XML pour analyser les chaînes XML. L'exemple suivant crée une table avec deux colonnes.
- La colonne
idest extraite de la chaîne JSON en tant que typeSTRING. variant_columncontient l'intégralité de la chaîne JSON encodée en tant que typeVARIANT.
CREATE TABLE table_name AS
SELECT json_string:id AS id,
PARSE_JSON(json_string) variant_column
FROM source_data
Databricks recommande d'extraire les champs fréquemment interrogés et de les stocker sous forme de colonnes non-variantes pour accélérer les requêtes et optimiser le layout du stockage.
VARIANT les colonnes ne peuvent pas être utilisées pour les clés de clustering, les partitions ou les clés de Z-order. Le type de données VARIANT ne peut pas être utilisé pour les comparaisons, les regroupements, le tri et les opérations d'ensemble. Pour plus d'informations, voir Limitations.
Insérer des données à l’aide de parse_json
Si la table cible contient déjà une colonne encodée en VARIANT, vous pouvez utiliser parse_json pour insérer des enregistrements de chaînes JSON en tant que VARIANT. Par exemple, analysez les chaînes JSON de la colonne json_string et insérez-les dans table_name.
- SQL
- Python
INSERT INTO table_name (variant_column)
SELECT PARSE_JSON(json_string)
FROM source_data
from pyspark.sql.functions import col, parse_json
(spark.read
.table("source_data")
.select(parse_json(col("json_string")))
.write
.mode("append")
.saveAsTable("table_name")
)
Insérer des données à l’aide de from_xml
Si la table cible contient déjà une colonne encodée en VARIANT, vous pouvez utiliser from_xml pour insérer des enregistrements de chaîne XML sous forme de VARIANT. Par exemple, analysez les chaînes XML de la colonne xml_string et insérez-les dans table_name.
- SQL
- Python
INSERT INTO table_name (variant_column)
SELECT FROM_XML(xml_string, 'variant')
FROM source_data
from pyspark.sql.functions import col, from_xml
(spark.read
.table("source_data")
.select(from_xml(col("xml_string"), "variant"))
.write
.mode("append")
.saveAsTable("table_name")
)
Insérer des données à l’aide de from_csv
Si la table cible contient déjà une colonne encodée en VARIANT, vous pouvez utiliser from_csv pour insérer des enregistrements de chaînes CSV en tant que VARIANT. Par exemple, analysez les enregistrements CSV de la colonne csv_string et insérez-les dans table_name.
- SQL
- Python
INSERT INTO table_name (variant_column)
SELECT FROM_CSV(csv_string, 'v variant').v
FROM source_data
from pyspark.sql.functions import col, from_csv
(spark.read
.table("source_data")
.select(from_csv(col("csv_string"), "v variant").v)
.write
.mode("append")
.saveAsTable("table_name")
)
Ingérer des données depuis le stockage d'objets cloud en tant que variante
Auto Loader peut être utilisé pour charger toutes les données des sources de fichiers prises en charge en tant que colonne VARIANT unique dans une table cible. Parce que VARIANT est flexible aux changements de schéma et de type et maintient la sensibilité à la casse et les valeurs NULL présentes dans la source de données, ce modèle est robuste à la plupart des scénarios d'ingestion, avec les mises en garde suivantes :
- Les enregistrements mal formés ne peuvent pas être encodés à l’aide du type
VARIANT. VARIANTle type ne peut contenir que des enregistrements d'une taille maximale de 16 Mo.
La variante traite les enregistrements trop volumineux de la même manière que les enregistrements corrompus. Dans le mode de traitement PERMISSIVE default, les enregistrements trop volumineux sont capturés dans le corruptRecordColumn.
Puisque l'enregistrement entier est enregistré en tant que colonne VARIANT unique, aucune évolution des schémas ne se produit pendant l'ingestion et rescuedDataColumn n'est pas pris en charge. L'exemple suivant suppose que la table cible existe déjà avec une seule colonne VARIANT.
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("singleVariantColumn", "variant_column")
.load("/Volumes/catalog_name/schema_name/volume_name/path")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("table_name")
)
Vous pouvez également spécifier VARIANT lors de la définition d'un schéma ou du passage de schemaHints. Les données du champ source référencé doivent contenir un enregistrement valide. Les exemples suivants illustrent cette syntaxe.
# Define the schema.
# Writes the columns `name` as a string and `address` as variant.
(spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.schema("name STRING, address VARIANT")
.load("/Volumes/catalog_name/schema_name/volume_name/path")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("table_name")
)
# Define the schema.
# A single field `payload` containing JSON data is written as variant.
(spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.schema("payload VARIANT")
.load("/Volumes/catalog_name/schema_name/volume_name/path")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("table_name")
)
# Supply schema hints.
# Writes the `address` column as variant.
# Infers the schema for other fields using standard rules.
(spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "address VARIANT")
.load("/Volumes/catalog_name/schema_name/volume_name/path")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("table_name")
)
Utiliser COPY INTO avec variante
Databricks recommande d'utiliser Auto Loader plutôt que COPY INTO lorsque disponible.
COPY INTO prend en charge l'ingestion de l'intégralité du contenu d'une source de données prise en charge en une seule colonne. L'exemple suivant crée une nouvelle table avec une seule colonne VARIANT, puis utilise COPY INTO pour ingérer des enregistrements à partir d'une source de fichiers JSON.
CREATE TABLE table_name (variant_column VARIANT);
COPY INTO table_name
FROM '/Volumes/catalog_name/schema_name/volume_name/path'
FILEFORMAT = JSON
FILES = ('file-name')
FORMAT_OPTIONS ('singleVariantColumn' = 'variant_column')
Stream Kafka data as variant
De nombreux streams Kafka encodent leurs charges utiles en JSON. L'ingestion de Kafka Stream à l'aide de VARIANT rend ces charges de travail robustes aux changements de schéma.
L'exemple suivant démontre la lecture d'une source de streaming Kafka, convertissant le key en STRING et le value en VARIANT, et l'écriture dans une table cible.
from pyspark.sql.functions import col, parse_json
(spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1")
.option("startingOffsets", "earliest")
.load()
.select(
col("key").cast("string"),
parse_json(col("value").cast("string"))
).writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("table_name")
)
Étapes suivantes
- Query des données de variante.
- Configurez la prise en charge du type de variante pour Apache Iceberg et Delta Lake.
- En savoir plus sur Auto Loader. Consultez Qu'est-ce qu'Auto Loader ?.