Aller au contenu principal

Charger des données dans les pipelines

Vous pouvez charger des données depuis n’importe quelle source de données prise en charge par Apache Spark sur Databricks à l’aide de pipelines. Vous pouvez définir des datasets (tables et vues) dans un pipeline pour toute query qui renvoie un DataFrame Spark, y compris les DataFrames de streaming et les DataFrames Pandas pour Spark. Pour les tâches d’ingestion de données, Databricks recommande l’utilisation de tables de streaming pour la plupart des cas d’usage. Les tables de streaming sont utiles pour ingérer des données depuis le stockage d’objets cloud à l’aide d’Auto Loader ou depuis des bus de messages comme Kafka. Pour en savoir plus sur les tables de streaming, le type de dataset principal pour l’ingestion, consultez Tables de streaming.

Toutes les sources de données ne prennent pas en charge SQL pour l'ingestion. Cependant, vous pouvez combiner des sources SQL et Python dans le même pipeline pour utiliser Python là où cela est nécessaire. Pour plus de détails sur l'utilisation de bibliothèques non packagées par default avec les pipelines, consultez Gérer les dépendances Python pour les pipelines. Pour des informations générales sur l'ingestion dans Databricks, consultez Choisir un connecteur standard.

Les exemples suivants illustrent certains modèles courants de chargement de données.

Identifiez vos sources de données et votre chemin de connexion

Avant d’écrire le code du pipeline, inventoriez toutes les origines de données. Pour chaque source, notez la manière dont elle expose ses données (fichiers, base de données, système SaaS, API ou stream), sa fréquence de modification, ainsi que les identifiants et l’accès réseau requis. La méthode de connexion détermine souvent si une source est naturellement de type batch ou streaming ; il est donc préférable de bien la définir dès le début pour éviter tout travail supplémentaire ultérieur.

Placez chaque source dans l’un des chemins de connexion suivants. Le tableau suivant liste le mécanisme préféré pour chaque source de pipeline :

Source

Chemin d'accès de la connexion

Fichiers arrivant dans le stockage d’objets cloud (S3, Azure Data Lake Storage, GCS)

Le point de départ le plus courant. Utilisez Auto Loader (format cloudFiles), qui gère la découverte incrémentielle, l’inférence de schéma et l’évolution des schémas. Voir Charger des fichiers à partir du stockage d’objets cloud.

Bases de données et applications SaaS (Salesforce, SQL Server, PostgreSQL, Workday)

Utilisez un connecteur géré Lakeflow Connect lorsqu'il en existe un pour votre source. Les connecteurs gérés sont pilotés par la configuration et gèrent l'authentification ainsi que l'extraction incrémentale ou CDC pour vous. Voir les concepts des connecteurs Lakeflow Connect. Si aucun connecteur géré n'existe pour votre source, ingérez-la directement ou enregistrez d'abord ses réponses sous forme de fichiers. Voir Ingérer des données à partir d'une API dans des pipelines.

Bus de messages (Kafka, Kinesis, Azure Event Hubs, Pub/Sub)

Lisez directement en tant que source Structured Streaming car il s’agit de sources de streaming natives. Voir Charger des données à partir d’un bus de messages.

Autres tables Delta ou assets Unity Catalog, y compris les tables produites par d'autres pipelines ou jobs

Référencez-les directement et laissez la gouvernance et la traçabilité de Unity Catalog gérer la découverte et l’accès. Voir Charger à partir d’une table existante.

Données de référence statiques ou de petite taille (fichiers de recherche, fichiers CSV changeant rarement)

Charger en tant que source batch dans une vue matérialisée. Le streaming d’éléments qui changent à peine n’offre aucun avantage. Voir Charger des datasets petits ou statiques à partir du stockage d’objets cloud.

Une API HTTP ou REST arbitraire sans connecteur géré

Effectuez une extraction depuis l'API dans le pipeline ou enregistrez d'abord ses réponses sous forme de fichiers. Voir Ingérer des données depuis une API dans des pipelines.

Source

Chemin d'accès de la connexion

Fichiers arrivant dans le stockage d’objets cloud (S3, Azure Data Lake Storage, GCS)

Le point de départ le plus courant. Utilisez Auto Loader (format cloudFiles), qui gère la découverte incrémentielle, l’inférence de schéma et l’évolution des schémas. Voir Charger des fichiers à partir du stockage d’objets cloud.

Bases de données et applications SaaS (Salesforce, SQL Server, PostgreSQL, Workday)

Utilisez un connecteur géré Lakeflow Connect lorsqu'il en existe un pour votre source. Les connecteurs gérés sont pilotés par la configuration et gèrent l'authentification ainsi que l'extraction incrémentale ou CDC pour vous. Voir les concepts des connecteurs Lakeflow Connect. Si aucun connecteur géré n'existe pour votre source, ingérez-la directement ou enregistrez d'abord ses réponses sous forme de fichiers. Voir Ingérer des données à partir d'une API dans des pipelines.

Bus de messages (Kafka, Kinesis, Azure Event Hubs, Pub/Sub)

Lisez directement en tant que source Structured Streaming car il s’agit de sources de streaming natives. Voir Charger des données à partir d’un bus de messages.

Autres tables Delta ou assets Unity Catalog, y compris les tables produites par d'autres pipelines ou jobs

Référencez-les directement et laissez la gouvernance et la traçabilité de Unity Catalog gérer la découverte et l’accès. Voir Charger à partir d’une table existante.

Données de référence statiques ou de petite taille (fichiers de recherche, fichiers CSV changeant rarement)

Charger en tant que source batch dans une vue matérialisée. Le streaming d’éléments qui changent à peine n’offre aucun avantage. Voir Charger des datasets petits ou statiques à partir du stockage d’objets cloud.

Une API HTTP ou REST arbitraire sans connecteur géré

Effectuez une extraction depuis l'API dans le pipeline ou enregistrez d'abord ses réponses sous forme de fichiers. Voir Ingérer des données depuis une API dans des pipelines.

Pour chaque source, confirmez les points suivants avant de procéder à la construction :

  • Identité : ce en tant que quoi le pipeline s'exécute. Les pipelines peuvent s'exécuter en tant que service principal ; configurez donc cela en premier pour éviter de dépendre d'un compte personnel.
  • Chemin réseau : connectivité requise par la source, telle qu’un identifiant de stockage, un emplacement externe ou la connectivité gérée par Lakeflow Connect.
  • Sémantique de modification : comment la source signale les mises à jour et les suppressions, le cas échéant. Cela détermine si vous avez besoin de la CDC ou si vous pouvez traiter la source comme étant en ajout seul.

Choisissez un format de fichier et une couche de stockage

Les pipelines prennent la majeure partie de ces décisions pour vous. Chaque table de streaming et vue matérialisée créée par un pipeline est stockée par default comme une table Delta, ce qui vous donne des transactions ACID, l'application des schémas et l'évolution, time travel, ainsi que la gouvernance et la lignée du Unity Catalog sur chaque jeu de données. Vous ne choisissez pas le format des sorties pipeline. Vos véritables décisions se situent aux deux extrémités du pipeline :

  • Format d’entrée brut : tout ce que la source produit, comme CSV, JSON ou Parquet. Auto Loader et read_files() les prennent en charge directement. Spécifiez le format avec cloudFiles.format en Python ou l’argument format => en SQL. Si vous contrôlez la source, préférez Parquet ou Avro, car ils intègrent le schéma et se compressent mieux, ce qui accélère l’ingestion et l’inférence de schéma. Les pipelines gèrent tous ces formats, ne laissez donc pas le format restreindre votre choix de source.
  • Emplacement de stockage brut : pour les fichiers, déposez les données dans un volume Unity Catalog plutôt que dans un chemin de compartiment non gouverné, afin que la traçabilité et le contrôle d'accès s'étendent jusqu'à la zone de dépôt. Voir Qu'est-ce que les volumes Unity Catalog ?.

Pour les tables produites par un pipeline, vos choix restants sont le catalogue cible et le schéma, qui définissent la frontière de gouvernance et la découvrabilité, ainsi que le Layout des grandes tables. Utilisez CLUSTER BY (cluster liquide) pour maintenir une bonne performance des requêtes à mesure que les tables grandissent sans ajuster manuellement les partitions. Voir Utiliser le regroupement liquide pour les tables.

Charger à partir d’une table existante

Chargez des données à partir de n'importe quelle table existante sur Databricks. Vous pouvez transformer les données à l'aide d'une query ou charger la table pour un traitement ultérieur dans votre pipeline.

Python
@dp.table(
comment="A table summarizing counts of the top baby names for New York for 2021."
)
def top_baby_names_2021():
return (
spark.read.table("baby_names_prepared")
.filter(expr("Year_Of_Birth == 2021"))
.groupBy("First_Name")
.agg(sum("Count").alias("Total_Count"))
.sort(desc("Total_Count"))
)

Charger des fichiers depuis le stockage d'objets cloud

Databricks recommande d'utiliser Auto Loader dans les pipelines pour la plupart des tâches d'ingestion de données à partir du stockage d'objets cloud ou à partir de fichiers dans un volume Unity Catalog. Auto Loader et les pipelines sont conçus pour charger de manière incrémentielle et idempotente les données en constante augmentation dès leur arrivée dans le stockage cloud. Voir Qu'est-ce qu'Auto Loader ? et Charger des données à partir du stockage d'objets.

L'exemple suivant lit des données à partir du stockage cloud à l'aide d'Auto Loader.

Python
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("s3://mybucket/analysis/*/*/*.json")
)

Les exemples suivants utilisent Auto Loader pour créer des datasets à partir de fichiers CSV dans un volume Unity Catalog.

Python
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/Volumes/my_catalog/retail_org/customers/")
)
remarque
  • Si vous utilisez Auto Loader avec des notifications de fichiers et exécutez un refresh complet pour votre pipeline ou votre table de streaming, vous devez nettoyer manuellement vos ressources. Vous pouvez utiliser le CloudFilesResourceManager dans un Notebook pour effectuer le nettoyage.
  • Pour charger des fichiers avec Auto Loader dans un pipeline activé pour Unity Catalog, vous devez utiliser des emplacements externes. Pour en savoir plus sur l'utilisation d'Unity Catalog avec des pipelines, consultez Utiliser Unity Catalog avec des pipelines.

S'authentifier auprès du stockage cloud

Auto Loader utilise les emplacements externes Unity Catalog pour s'authentifier auprès du stockage cloud. Vous devez configurer un emplacement externe pour le chemin de stockage à partir duquel vous souhaitez lire et accorder le privilège READ FILES à l'utilisateur exécutant.

Pour ingérer depuis Amazon S3, configurez un emplacement externe adossé à un identifiant de stockage qui fait référence à un bucket S3. Pour plus d'informations, consultez Se connecter au stockage d'objets cloud à l'aide de Unity Catalog.

Charger les données depuis un bus de messages

Vous pouvez configurer les pipelines pour ingérer des données à partir de bus de messages. Databricks recommande d'utiliser des tables de streaming avec exécution continue et autoscaling amélioré pour offrir l'ingestion la plus efficace pour le chargement à faible latence à partir de bus de messages. Pour plus d'informations, consultez Optimisez l'utilisation des clusters LakeFlow Pipelines avec l'autoscaling.

Par exemple, le code suivant configure une table de streaming pour ingérer des données de Kafka en utilisant la fonction read_kafka.

Python
from pyspark import pipelines as dp

@dp.table
def kafka_raw():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka_server:9092")
.option("subscribe", "topic1")
.load()
)

Ingérer à partir de Google Pub/Sub

L'exemple suivant crée une table de streaming qui lit à partir d'un sujet Google Pub/Sub à l'aide de la fonction read_pubsub.

Python
@dp.table
def pubsub_raw():
auth_options = {
"clientId": client_id,
"clientEmail": client_email,
"privateKey": private_key,
"privateKeyId": private_key_id
}
return (
spark.readStream
.format("pubsub")
.option("subscriptionId", "my-subscription")
.option("topicId", "my-topic")
.option("projectId", "my-project")
.options(auth_options)
.load()
)

Databricks recommande d'utiliser des secrets lors de la fourniture d'options d'autorisation. Consultez Configurer l'accès à Pub/Sub pour toutes les options d'authentification.

Pour ingérer à partir d'autres sources de bus de messages :

Charger des données à partir d’Azure Event Hubs

Azure Event Hubs est un service de streaming de données qui fournit une interface compatible avec Apache Kafka. Vous pouvez utiliser le connecteur Kafka de Structured Streaming, inclus dans l'exécution du pipeline, pour charger des messages depuis Azure Event Hub. Pour en savoir plus sur le chargement et le traitement des messages depuis Azure Event Hubs, consultez Utiliser Azure Event Hubs comme source de données de pipeline.

Charger des données à partir de systèmes externes

Les pipelines prennent en charge le chargement de données à partir de n'importe quelle source de données prise en charge par Databricks. Voir Connecter aux sources de données et aux services externes. Vous pouvez également charger des données externes à l'aide de Lakehouse Federation pour les sources de données prises en charge. Étant donné que Lakehouse Federation nécessite Databricks Runtime 13.3 LTS ou supérieur, pour utiliser Lakehouse Federation, configurez votre pipeline pour utiliser le canal de distribution en avant-première.

Certaines sources de données n’ont pas de prise en charge SQL équivalente. Si vous ne pouvez pas utiliser la Lakehouse Federation avec l'une de ces sources de données, vous pouvez utiliser Python pour ingérer des données depuis la source. Vous pouvez ajouter des fichiers sources Python et SQL au même pipeline. L’exemple suivant déclare une vue matérialisée pour accéder à l’état actuel des données dans une table PostgreSQL distante.

Python
import dp

@dp.table
def postgres_raw():
return (
spark.read
.format("postgresql")
.option("dbtable", table_name)
.option("host", database_host_url)
.option("port", 5432)
.option("database", database_name)
.option("user", username)
.option("password", password)
.load()
)

Charger de petits datasets ou des datasets statiques à partir du stockage d'objets cloud

Vous pouvez charger des datasets petits ou statiques à l'aide de la syntaxe de chargement d'Apache Spark. Les pipelines prennent en charge tous les formats de fichier pris en charge par Apache Spark sur Databricks. Pour obtenir la liste complète, consultez les options de format de données.

Les exemples suivants montrent comment charger des JSON pour créer une table.

Python
@dp.table
def clickstream_raw():
return (spark.read.format("json").load("/databricks-datasets/wikipedia-datasets/data-001/clickstream/raw-uncompressed-json/2015_2_clickstream.json"))
remarque

La fonction SQL read_files est commune à tous les environnements SQL sur Databricks. C'est le modèle recommandé pour l'accès direct aux fichiers à l'aide de SQL dans les pipelines. Pour plus d'informations, consultez Options.

Charger des données à partir d'une source de données Python personnalisée

Les sources de données personnalisées Python vous permettent de charger des données dans des formats personnalisés. Vous pouvez écrire du code pour lire à partir d'une source de données externe spécifique et y écrire, ou utiliser votre code Python existant pour lire les données de vos propres systèmes internes. Pour plus de détails sur le développement de sources de données Python, consultez sources de données personnalisées PySpark.

L'exemple suivant enregistre une source de données personnalisée avec le nom de format my_custom_datasource et en lit les données en modes batch et streaming.

Python
from pyspark import pipelines as dp

# Assume `my_custom_datasource` is a custom Python custom data
# source that supports both batch and streaming reads, and has
# been registered using `spark.dataSource.register`.

# This creates a materialized view
@dp.table(name = "read_from_batch")
def read_from_batch():
return spark.read.format("my_custom_datasource").load()

# This creates a streaming table
@dp.table(name = "read_from_streaming")
def read_from_streaming():
return spark.readStream.format("my_custom_datasource").load()

Configurez une table de streaming pour ignorer les modifications dans une table de streaming source

Par default, les tables de streaming nécessitent des sources en mode ajout uniquement. Si votre table de streaming source nécessite des mises à jour ou des suppressions (par exemple, pour le traitement du « droit à l'oubli » du GDPR), utilisez l'indicateur skipChangeCommits pour ignorer ces changements. Ce flag ne fonctionne qu'avec spark.readStream à l'aide de la fonction option() et ne peut pas être utilisé lorsque la table de streaming source est la cible d'une fonction create_auto_cdc_flow(). Pour plus d'informations, consultez Gérer les modifications apportées aux tables Delta Lake sources.

Python
@dp.table
def b():
return spark.readStream.option("skipChangeCommits", "true").table("A")

Accédez en toute sécurité aux identifiants de stockage avec des secrets dans un pipeline

Vous pouvez utiliser les secrets Databricks pour stocker des identifiants tels que des clés d'accès ou des mots de passe. Pour configurer le secret dans votre pipeline, utilisez une propriété Spark dans la configuration du cluster des paramètres du pipeline. Consultez Configurer le compute classique pour les pipelines.

L’exemple suivant utilise un secret pour stocker une clé d’accès requise pour lire les données d’entrée d’un compte de stockage Azure Data Lake Storage à l’aide d’Auto Loader. Vous pouvez utiliser cette même méthode pour configurer tout secret requis par votre pipeline, par exemple, les clés AWS pour accéder à S3, ou le mot de passe d’un Hive metastore Apache.

Pour en savoir plus sur l'utilisation d'Azure Data Lake Storage, consultez Connecter à Azure Data Lake Storage et au Stockage Blob.

remarque

Vous devez ajouter le préfixe spark.hadoop. à la clé de configuration spark_conf qui définit la valeur secrète.

JSON
{
"id": "43246596-a63f-11ec-b909-0242ac120002",
"storage": "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/<path>",
"clusters": [
{
"spark_conf": {
"spark.hadoop.fs.azure.account.key.<storage-account-name>.dfs.core.windows.net": "{{secrets/<scope-name>/<secret-name>}}"
},
"autoscale": {
"min_workers": 1,
"max_workers": 5,
"mode": "ENHANCED"
}
}
],
"development": true,
"continuous": false,
"libraries": [
{
"notebook": {
"path": "/Users/user@databricks.com/Pipeline Notebooks/pipeline quickstart"
}
}
],
"name": "pipeline quickstart using ADLS2"
}

Dans cet exemple de code, remplacez les valeurs suivantes.

Espace réservé

Remplacer par

<container-name>

Le nom du conteneur de compte de stockage Azure.

<storage-account-name>

Le nom du compte de stockage ADLS.

<path>

Le chemin d'accès pour les données de sortie et les métadonnées du pipeline.

<scope-name>

Le nom du Secret Scope Databricks.

<secret-name>

Nom de la clé contenant la clé d'accès du compte de stockage Azure.

Espace réservé

Remplacer par

<container-name>

Le nom du conteneur de compte de stockage Azure.

<storage-account-name>

Le nom du compte de stockage ADLS.

<path>

Le chemin d'accès pour les données de sortie et les métadonnées du pipeline.

<scope-name>

Le nom du Secret Scope Databricks.

<secret-name>

Nom de la clé contenant la clé d'accès du compte de stockage Azure.

Python
from pyspark import pipelines as dp

json_path = "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/<path-to-input-dataset>"
@dp.create_table(
comment="Data ingested from an ADLS2 storage account."
)
def read_from_ADLS2():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load(json_path)
)

Dans cet exemple de code, remplacez les valeurs suivantes.

Espace réservé

Remplacer par

<container-name>

Le nom du conteneur de compte de stockage Azure qui stocke les données d'entrée.

<storage-account-name>

Le nom du compte de stockage ADLS.

<path-to-input-dataset>

Le chemin d'accès au dataset d'entrée.

Espace réservé

Remplacer par

<container-name>

Le nom du conteneur de compte de stockage Azure qui stocke les données d'entrée.

<storage-account-name>

Le nom du compte de stockage ADLS.

<path-to-input-dataset>

Le chemin d'accès au dataset d'entrée.