Aller au contenu principal

Tutoriel : COPY INTO avec Spark SQL

Databricks vous recommande d’utiliser la commande COPY INTO pour le chargement de données incrémentiel et en bloc pour les sources de données contenant des milliers de fichiers.

Dans ce tutoriel, vous utilisez la commande COPY INTO pour charger des données JSON depuis un volume Unity Catalog dans une table Delta de votre workspace Databricks. Vous utilisez le dataset échantillon Wanderbricks comme source de données. Pour les cas d'utilisation d'ingestion plus avancés, consultez What is Auto Loader?.

Exigences

Étape 1 : Configurez votre environnement

Le code de ce tutoriel utilise un volume Unity Catalog pour stocker les fichiers source JSON. Remplacez <catalog> par un catalogue où vous disposez des autorisations CREATE SCHEMA et CREATE VOLUME. Si vous ne pouvez pas exécuter le code, contactez l'administrateur de votre workspace.

Créez un notebook et associez-le à une ressource de compute. Exécutez ensuite le code suivant pour configurer un schéma et un volume pour ce didacticiel.

Python
# Set parameters and reset demo environment

catalog = "<catalog>"

username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first()[0]
schema = f"copyinto_{username}_db"
volume = "copy_into_source"
source = f"/Volumes/{catalog}/{schema}/{volume}"

spark.sql(f"SET c.catalog={catalog}")
spark.sql(f"SET c.schema={schema}")
spark.sql(f"SET c.volume={volume}")

spark.sql(f"DROP SCHEMA IF EXISTS {catalog}.{schema} CASCADE")
spark.sql(f"CREATE SCHEMA {catalog}.{schema}")
spark.sql(f"CREATE VOLUME {catalog}.{schema}.{volume}")

Étape 2 : Écrivez des exemples de données dans le volume au format JSON

La commande COPY INTO charge les données à partir de sources basées sur des fichiers. Lisez depuis la table d'exemple Wanderbricks bookings et écrivez un **batch** d'enregistrements sous forme de fichiers JSON vers votre volume, en simulant l'arrivée de données d'un système externe.

Python
# Write a batch of Wanderbricks bookings data as JSON to the volume

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_1 = bookings.orderBy("booking_id").limit(20)
batch_1.write.mode("append").json(f"{source}/bookings")

Étape 3 : utiliser COPY INTO pour charger des données JSON de manière idempotente

Créer une table Delta cible avant d'utiliser COPY INTO. Vous n'avez rien d'autre à fournir qu'un nom de table dans votre instruction CREATE TABLE. Puisque cette action est idempotente, Databricks charge les données une seule fois, même si vous exécutez le code plusieurs fois.

Python
# Create target table and load data

spark.sql(f"CREATE TABLE IF NOT EXISTS {catalog}.{schema}.bookings_target")

spark.sql(f"""
COPY INTO {catalog}.{schema}.bookings_target
FROM '/Volumes/{catalog}/{schema}/{volume}/bookings'
FILEFORMAT = JSON
FORMAT_OPTIONS ('mergeSchema' = 'true')
COPY_OPTIONS ('mergeSchema' = 'true')
""")

Étape 4  : Prévisualiser le contenu de votre table

Vérifiez que la table contient 20 lignes du premier lot de données de réservations Wanderbricks et que le schéma a été correctement inféré à partir des fichiers source JSON.

Python
# Review loaded data

display(spark.sql(f"SELECT * FROM {catalog}.{schema}.bookings_target"))

Étape 5 : Chargez plus de données et prévisualisez les résultats

Vous pouvez simuler l'arrivée de données supplémentaires provenant d'un système externe en écrivant un autre lot d'enregistrements et en exécutant COPY INTO à nouveau. Exécutez le code suivant pour écrire un second batch de données.

Python
# Write another batch of Wanderbricks bookings data as JSON

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_2 = bookings.orderBy(bookings.booking_id.desc()).limit(20)
batch_2.write.mode("append").json(f"{source}/bookings")

Exécutez ensuite la commande COPY INTO de l'étape 3 à nouveau et prévisualisez la table pour confirmer les nouveaux enregistrements. Seuls les nouveaux fichiers sont chargés.

Python
# Confirm new data was loaded

display(spark.sql(f"SELECT COUNT(*) AS total_rows FROM {catalog}.{schema}.bookings_target"))

Étape 6 : Nettoyer le tutoriel

Lorsque vous avez terminé ce tutoriel, vous pouvez nettoyer les Ressources associées si vous ne souhaitez plus les conserver. Supprimez le schéma, les tables et le volume, et supprimez toutes les données.

Python
# Drop schema and all associated objects

spark.sql(f"DROP SCHEMA IF EXISTS {catalog}.{schema} CASCADE")

Ressources supplémentaires