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
- Accès à une ressource de compute. See compute.
- Un Workspace compatible Unity Catalog avec les autorisations pour créer des schémas et des volumes dans un catalogue. Voir Se connecter au stockage d'objets cloud à l'aide de Unity Catalog.
É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
- SQL
# 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}")
-- Reset demo environment
DROP SCHEMA IF EXISTS <catalog>.copy_into_tutorial CASCADE;
CREATE SCHEMA <catalog>.copy_into_tutorial;
CREATE VOLUME <catalog>.copy_into_tutorial.copy_into_source;
É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
- SQL
# 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")
L'écriture de fichiers sur un volume nécessite Python. Dans un workflow réel, ces données proviendraient 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("/Volumes/<catalog>/copy_into_tutorial/copy_into_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
- SQL
# 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')
""")
-- Create target table and load data
CREATE TABLE IF NOT EXISTS <catalog>.copy_into_tutorial.bookings_target;
COPY INTO <catalog>.copy_into_tutorial.bookings_target
FROM '/Volumes/<catalog>/copy_into_tutorial/copy_into_source/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
- SQL
# Review loaded data
display(spark.sql(f"SELECT * FROM {catalog}.{schema}.bookings_target"))
-- Review loaded data
SELECT * FROM <catalog>.copy_into_tutorial.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
- SQL
# 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")
L'écriture de fichiers sur un volume nécessite Python. Dans un workflow réel, ces données proviendraient d'un système externe.
%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("/Volumes/<catalog>/copy_into_tutorial/copy_into_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
- SQL
# Confirm new data was loaded
display(spark.sql(f"SELECT COUNT(*) AS total_rows FROM {catalog}.{schema}.bookings_target"))
-- Confirm new data was loaded
SELECT COUNT(*) AS total_rows FROM <catalog>.copy_into_tutorial.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
- SQL
# Drop schema and all associated objects
spark.sql(f"DROP SCHEMA IF EXISTS {catalog}.{schema} CASCADE")
-- Drop schema and all associated objects
DROP SCHEMA IF EXISTS <catalog>.copy_into_tutorial CASCADE;