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 mélanger des sources SQL et Python dans le même pipeline pour utiliser Python là où cela est nécessaire. Pour en savoir plus sur l'utilisation des bibliothèques non fournies 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 les connecteurs standard dans Lakeflow Connect.
Les exemples suivants illustrent certains modèles courants de chargement de données.
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
- SQL
@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"))
)
CREATE OR REFRESH MATERIALIZED VIEW top_baby_names_2021
COMMENT "A table summarizing counts of the top baby names for New York for 2021."
AS SELECT
First_Name,
SUM(Count) AS Total_Count
FROM baby_names_prepared
WHERE Year_Of_Birth = 2021
GROUP BY First_Name
ORDER BY Total_Count DESC
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
- SQL
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("s3://mybucket/analysis/*/*/*.json")
)
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT *
FROM STREAM read_files(
's3://mybucket/analysis/*/*/*.json',
format => "json"
);
Les exemples suivants utilisent Auto Loader pour créer des datasets à partir de fichiers CSV dans un volume Unity Catalog.
- Python
- SQL
@dp.table
def customers():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/Volumes/my_catalog/retail_org/customers/")
)
CREATE OR REFRESH STREAMING TABLE customers
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/retail_org/customers/",
format => "csv"
)
- 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
- SQL
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()
)
CREATE OR REFRESH STREAMING TABLE kafka_raw AS
SELECT *
FROM STREAM read_kafka(
bootstrapServers => 'kafka_server:9092',
subscribe => 'topic1'
);
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
- SQL
@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()
)
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'my-subscription',
projectId => 'my-project',
topicId => 'my-topic',
clientEmail => secret('pubsub-scope', 'clientEmail'),
clientId => secret('pubsub-scope', 'clientId'),
privateKeyId => secret('pubsub-scope', 'privateKeyId'),
privateKey => secret('pubsub-scope', 'privateKey')
);
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 :
- Kinesis : read_kinesis
- Pulsar : read_pulsar
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.
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
- SQL
@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"))
CREATE OR REFRESH MATERIALIZED VIEW clickstream_raw
AS SELECT * FROM read_files(
"/databricks-datasets/wikipedia-datasets/data-001/clickstream/raw-uncompressed-json/2015_2_clickstream.json"
)
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.
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.
@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.
Vous devez ajouter le préfixe spark.hadoop. à la clé de configuration spark_conf qui définit la valeur secrète.
{
"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 |
|---|---|
| Le nom du conteneur de compte de stockage Azure. |
| Le nom du compte de stockage ADLS. |
| Le chemin d'accès pour les données de sortie et les métadonnées du pipeline. |
| Le nom du Secret Scope Databricks. |
| Nom de la clé contenant la clé d'accès du compte de stockage Azure. |
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 |
|---|---|
| Le nom du conteneur de compte de stockage Azure qui stocke les données d'entrée. |
| Le nom du compte de stockage ADLS. |
| Le chemin d'accès au dataset d'entrée. |