Commencez à utiliser COPY INTO pour charger des données
Vous pouvez utiliser la commande SQL COPY INTO pour charger des données d'un emplacement de fichier dans une table Delta. COPY INTO est relançable et idempotent — les fichiers à l'emplacement source qui ont déjà été chargés sont ignorés lors des exécutions ultérieures.
COPY INTO offre ces fonctionnalités :
- Filtres de fichiers ou de dossiers facilement configurables à partir du stockage cloud, y compris les volumes S3, ADLS, ABFS, GCS et Unity Catalog.
- Prise en charge de plusieurs formats de fichiers source : CSV, JSON, XML, Avro, ORC, Parquet, texte et fichiers binaires.
- Traitement de fichiers exactement une fois (idempotent) par default.
- Inférence, mappage, fusion et évolution du schéma de la table cible.
Pour une expérience d'ingestion de fichiers plus évolutive et robuste, Databricks recommande aux utilisateurs SQL d'utiliser les tables de streaming. Pour plus d'informations, consultez Tables de streaming.
COPY INTO respecte le paramètre de Workspace pour les vecteurs de suppression. S'ils sont activés, les vecteurs de suppression sont activés sur la table cible lorsque COPY INTO s'exécute sur un SQL warehouse ou un compute exécutant Databricks Runtime 14.0 ou une version supérieure. Une fois les vecteurs de suppression activés, ils bloquent les query sur une table dans Databricks Runtime 11.3 LTS et les versions antérieures. Voir Vecteurs de suppression dans Databricks et Activer automatiquement les vecteurs de suppression.
Avant de commencer
Un administrateur de compte doit suivre les étapes de Configuration de l'accès aux données pour l'ingestion pour configurer l'accès aux données dans le stockage d'objets cloud avant que les utilisateurs ne puissent charger des données à l'aide de COPY INTO.
Charger les données dans une table Delta Lake sans schéma
Dans Databricks Runtime 11.3 LTS et versions ultérieures, vous pouvez créer des tables Delta de type espace réservé vides afin que le schéma soit inféré lors d'une commande COPY INTO en définissant mergeSchema sur true dans COPY_OPTIONS. L'exemple suivant utilise le Wanderbricks dataset. Remplacez <catalog>, <schema>, et <volume> par un catalogue, un schéma et un volume où vous avez les autorisations CREATE TABLE.
- SQL
- Python
- R
- Scala
CREATE TABLE IF NOT EXISTS <catalog>.<schema>.booking_updates_schemaless;
COPY INTO <catalog>.<schema>.booking_updates_schemaless
FROM '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
FILEFORMAT = JSON
FORMAT_OPTIONS ('mergeSchema' = 'true', 'multiLine' = 'true')
COPY_OPTIONS ('mergeSchema' = 'true');
table_name = '<catalog>.<schema>.booking_updates_schemaless'
source_data = '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
source_format = 'JSON'
spark.sql("CREATE TABLE IF NOT EXISTS " + table_name)
spark.sql("COPY INTO " + table_name + \
" FROM '" + source_data + "'" + \
" FILEFORMAT = " + source_format + \
" FORMAT_OPTIONS ('mergeSchema' = 'true', 'multiLine' = 'true')" + \
" COPY_OPTIONS ('mergeSchema' = 'true')"
)
library(SparkR)
sparkR.session()
table_name = "<catalog>.<schema>.booking_updates_schemaless"
source_data = "/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates"
source_format = "JSON"
sql(paste("CREATE TABLE IF NOT EXISTS ", table_name, sep = ""))
sql(paste("COPY INTO ", table_name,
" FROM '", source_data, "'",
" FILEFORMAT = ", source_format,
" FORMAT_OPTIONS ('mergeSchema' = 'true', 'multiLine' = 'true')",
" COPY_OPTIONS ('mergeSchema' = 'true')",
sep = ""
))
val table_name = "<catalog>.<schema>.booking_updates_schemaless"
val source_data = "/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates"
val source_format = "JSON"
spark.sql("CREATE TABLE IF NOT EXISTS " + table_name)
spark.sql("COPY INTO " + table_name +
" FROM '" + source_data + "'" +
" FILEFORMAT = " + source_format +
" FORMAT_OPTIONS ('mergeSchema' = 'true', 'multiLine' = 'true')" +
" COPY_OPTIONS ('mergeSchema' = 'true')"
)
Cette instruction SQL est idempotente. Cela signifie que vous pouvez le programmer pour qu'il s'exécute de manière répétée, et qu'il ne chargera de nouvelles données que dans votre table Delta.
La table Delta vide n'est pas utilisable en dehors de COPY INTO. INSERT INTO et MERGE INTO ne sont pas pris en charge pour écrire des données dans les tables Delta sans schéma. Une fois les données insérées dans la table avec COPY INTO, la table devient interrogeable.
Voir Créer des tables cibles pour COPY INTO.
Définir le schéma et charger les données dans une table Delta Lake
L'exemple suivant crée une table Delta et utilise la commande SQL COPY INTO pour charger des exemples de données du jeu de données Wanderbricks dans la table. Les fichiers source sont des fichiers JSON stockés dans un volume Unity Catalog. Vous pouvez exécuter l'exemple de code Python, R, Scala ou SQL depuis un Notebook attaché à un cluster Databricks. Vous pouvez également exécuter le code SQL à partir d'une query associée à un SQL Warehouse dans Databricks SQL. Remplacez <catalog>, <schema> et <volume> par un catalogue, un schéma et un volume où vous disposez des autorisations CREATE TABLE.
- SQL
- Python
- R
- Scala
DROP TABLE IF EXISTS <catalog>.<schema>.booking_updates_upload;
CREATE TABLE <catalog>.<schema>.booking_updates_upload (
booking_id BIGINT,
user_id BIGINT,
status STRING,
total_amount DOUBLE
);
COPY INTO <catalog>.<schema>.booking_updates_upload
FROM '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
FILEFORMAT = JSON
FORMAT_OPTIONS ('multiLine' = 'true');
SELECT * FROM <catalog>.<schema>.booking_updates_upload;
table_name = '<catalog>.<schema>.booking_updates_upload'
source_data = '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
source_format = 'JSON'
spark.sql("DROP TABLE IF EXISTS " + table_name)
spark.sql("CREATE TABLE " + table_name + " (" \
"booking_id BIGINT, " + \
"user_id BIGINT, " + \
"status STRING, " + \
"total_amount DOUBLE)"
)
spark.sql("COPY INTO " + table_name + \
" FROM '" + source_data + "'" + \
" FILEFORMAT = " + source_format + \
" FORMAT_OPTIONS ('multiLine' = 'true')"
)
booking_updates_upload_data = spark.sql("SELECT * FROM " + table_name)
display(booking_updates_upload_data)
library(SparkR)
sparkR.session()
table_name = "<catalog>.<schema>.booking_updates_upload"
source_data = "/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates"
source_format = "JSON"
sql(paste("DROP TABLE IF EXISTS ", table_name, sep = ""))
sql(paste("CREATE TABLE ", table_name, " (",
"booking_id BIGINT, ",
"user_id BIGINT, ",
"status STRING, ",
"total_amount DOUBLE)",
sep = ""
))
sql(paste("COPY INTO ", table_name,
" FROM '", source_data, "'",
" FILEFORMAT = ", source_format,
" FORMAT_OPTIONS ('multiLine' = 'true')",
sep = ""
))
booking_updates_upload_data = tableToDF(table_name)
display(booking_updates_upload_data)
val table_name = "<catalog>.<schema>.booking_updates_upload"
val source_data = "/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates"
val source_format = "JSON"
spark.sql("DROP TABLE IF EXISTS " + table_name)
spark.sql("CREATE TABLE " + table_name + " (" +
"booking_id BIGINT, " +
"user_id BIGINT, " +
"status STRING, " +
"total_amount DOUBLE)"
)
spark.sql("COPY INTO " + table_name +
" FROM '" + source_data + "'" +
" FILEFORMAT = " + source_format +
" FORMAT_OPTIONS ('multiLine' = 'true')"
)
val booking_updates_upload_data = spark.table(table_name)
display(booking_updates_upload_data)
Pour nettoyer, exécutez le code suivant pour supprimer la table d'exemple.
- SQL
- Python
- R
- Scala
DROP TABLE <catalog>.<schema>.booking_updates_upload
spark.sql("DROP TABLE " + table_name)
sql(paste("DROP TABLE ", table_name, sep = ""))
spark.sql("DROP TABLE " + table_name)
Nettoyage des fichiers de métadonnées
Vous pouvez exécuter VACUUM pour nettoyer les fichiers de métadonnées non référencés créés par COPY INTO dans Databricks Runtime 15.2 et versions supérieures.
Ressources supplémentaires
-
Modèles courants de chargement de données utilisant
COPY INTO. -
Databricks Runtime 7.x ou une version ultérieure :
COPY INTO