Aller au contenu principal

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

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.

SQL
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 id est extraite de la chaîne JSON en tant que type STRING.
  • variant_column contient l'intégralité de la chaîne JSON encodée en tant que type VARIANT.
SQL
CREATE TABLE table_name AS
SELECT json_string:id AS id,
PARSE_JSON(json_string) variant_column
FROM source_data
remarque

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
INSERT INTO table_name (variant_column)
SELECT PARSE_JSON(json_string)
FROM source_data

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
INSERT INTO table_name (variant_column)
SELECT FROM_XML(xml_string, 'variant')
FROM source_data

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
INSERT INTO table_name (variant_column)
SELECT FROM_CSV(csv_string, 'v variant').v
FROM source_data

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.
  • VARIANT le type ne peut contenir que des enregistrements d'une taille maximale de 16 Mo.
remarque

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.

Python
(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.

Python
# 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.

SQL
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.

Python
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