Aller au contenu principal

Référence linguistique Python des Lakeflow Pipelines

L'interface Python des Lakeflow pipelines est définie dans le module pyspark.pipelines, importé comme dp.

pipelines présentation du module​

Les fonctions Python des LakeFlow pipelines sont définies dans le module pyspark.pipelines (importées sous le nom dp). Vos pipelines implémentés avec l'API Python doivent importer ce module :

Python
from pyspark import pipelines as dp
remarque

La majeure partie du module de pipeline est uniquement disponible dans le contexte d'un pipeline. Pour plus d'informations sur l'édition du code de pipeline, consultez Develop and debug ETL pipelines with the Lakeflow Pipelines Editor.

Vous pouvez également appliquer @materialized_view et @table en dehors d'un pipeline, dans un notebook sur le compute général serverless, pour définir une vue matérialisée ou une table de streaming autonome. Consultez Définir des tables avec les décorateurs de pipelines.

pipelines Apache Spark™​

Apache Spark inclut des pipelines déclaratifs à partir de Spark 4.1, disponibles via le module pyspark.pipelines. Le Databricks Runtime étend ces capacités open source avec des APIs et des intégrations supplémentaires pour une utilisation en production gérée.

Le code écrit avec le module open source pipelines s'exécute sans modification sur Databricks. Les fonctionnalités suivantes ne font pas partie d'Apache Spark :

  • dp.create_auto_cdc_flow
  • dp.create_auto_cdc_from_snapshot_flow
  • @dp.expect(...)

Le module pipelines s'appelait auparavant dlt dans Databricks. Pour plus de détails, et davantage d'informations sur les différences avec Apache Spark, consultez Qu'est-il arrivé à @dlt?.

Fonctions pour les définitions de datasets​

Les pipeline utilisent des décorateurs Python pour définir des dataset tels que les vues matérialisées et les tables de streaming. Voir Fonctions pour définir des datasets.

Référence de l'API​

Exigences de codage pour les pipelines Python​

Voici les exigences importantes lorsque vous implémentez des pipelines avec l'interface Python des Lakeflow Pipelines :

  • Lakeflow pipelines évaluent le code qui définit un pipeline plusieurs fois pendant la planification et les exécutions de pipelines. Les fonctions Python qui définissent les datasets doivent inclure uniquement le code requis pour définir la table ou la vue. Une logique Python arbitraire incluse dans les définitions de dataset peut entraîner un comportement inattendu.
  • N'essayez pas d'implémenter une logique de monitoring personnalisée dans vos définitions de dataset. Voir Définir le monitoring personnalisé des pipelines avec des hooks d'événements.
  • La fonction utilisée pour définir un dataset doit renvoyer un DataFrame Spark. N'incluez pas de logique dans vos définitions de dataset qui n'est pas liée à un DataFrame renvoyé.
  • N'utilisez jamais de méthodes qui enregistrent ou écrivent dans des fichiers ou des tables dans le cadre de votre code de dataset de pipeline.

Exemples d'Opérations Apache Spark qui ne devraient jamais être utilisées dans le code de pipeline :

  • collect()
  • count()
  • toPandas()
  • save()
  • saveAsTable()
  • start()
  • toTable()

Qu'est-il arrivé à @dlt?​

Auparavant, Databricks utilisait le module dlt pour prendre en charge la fonctionnalité de pipeline. Le module dlt a été remplacé par le module pyspark.pipelines. Vous pouvez toujours utiliser dlt, mais Databricks recommande d'utiliser pipelines.

Différences entre DLT, Lakeflow Pipelines et Apache Spark Declarative Pipelines​

Le tableau suivant montre les différences de syntaxe et de fonctionnalité entre les DLT, les LakeFlow Pipelines et les Apache Spark Declarative Pipelines.

Pour une comparaison au niveau des fonctionnalités de ce que les Lakeflow Pipelines partagent avec et ajoutent à Apache Spark Declarative Pipelines, consultez Apache Spark Declarative Pipelines.

Pour un mappage propriété par propriété de la configuration de pipeline à la spécification de projet SDP, consultez la référence des propriétés de pipeline.

remarque

Dans la documentation Databricks, le produit Databricks est appelé LakeFlow Pipelines , et le framework open source qu'il étend est Apache Spark™ Declarative Pipelines ( SDP ). Les deux sont interopérables, mais diffèrent par leurs fonctionnalités — par exemple, les AUTO CDC APIs sont disponibles uniquement dans les pipelines Lakeflow.

Zone (Area)

Syntaxe DLT

Syntaxe SDP (Lakeflow et Apache, le cas échéant)

Disponible dans Apache Spark

Importations

import dlt

from pyspark import pipelines (as dp, facultatif)

Oui

Table de streaming

@dlt.table avec une trame de données en streaming

@dp.table

Oui

Vue matérialisée

@dlt.table avec un dataframe batch

@dp.materialized_view

Oui

Afficher

@dlt.view

@dp.temporary_view

Oui

Flux d'ajout

@dlt.append_flow

@dp.append_flow

Oui

Mettre à jour le flux

Indisponible

@dp.update_flow

Non

SQL – streaming

CREATE STREAMING TABLE ...

CREATE STREAMING TABLE ...

Oui

SQL – matérialisées

CREATE MATERIALIZED VIEW ...

CREATE MATERIALIZED VIEW ...

Oui

SQL – flux

CREATE FLOW ...

CREATE FLOW ...

Oui

Journal des événements

spark.read.table("event_log")

spark.read.table("event_log")

Non

Appliquer les modifications (CDC)

dlt.apply_changes(...)

dp.create_auto_cdc_flow(...)

Non

Attentes

@dlt.expect(...)

dp.expect(...)

Non

Mode continu

Configuration de pipeline avec Trigger continu

(même)

Non

Puits

@dlt.create_sink(...)

dp.create_sink(...)

Oui

ForEachBatch récepteur

Indisponible

@dp.foreach_batch_sink(...)

Non

Zone (Area)

Syntaxe DLT

Syntaxe SDP (Lakeflow et Apache, le cas échéant)

Disponible dans Apache Spark

Importations

import dlt

from pyspark import pipelines (as dp, facultatif)

Oui

Table de streaming

@dlt.table avec une trame de données en streaming

@dp.table

Oui

Vue matérialisée

@dlt.table avec un dataframe batch

@dp.materialized_view

Oui

Afficher

@dlt.view

@dp.temporary_view

Oui

Flux d'ajout

@dlt.append_flow

@dp.append_flow

Oui

Mettre à jour le flux

Indisponible

@dp.update_flow

Non

SQL – streaming

CREATE STREAMING TABLE ...

CREATE STREAMING TABLE ...

Oui

SQL – matérialisées

CREATE MATERIALIZED VIEW ...

CREATE MATERIALIZED VIEW ...

Oui

SQL – flux

CREATE FLOW ...

CREATE FLOW ...

Oui

Journal des événements

spark.read.table("event_log")

spark.read.table("event_log")

Non

Appliquer les modifications (CDC)

dlt.apply_changes(...)

dp.create_auto_cdc_flow(...)

Non

Attentes

@dlt.expect(...)

dp.expect(...)

Non

Mode continu

Configuration de pipeline avec Trigger continu

(même)

Non

Puits

@dlt.create_sink(...)

dp.create_sink(...)

Oui

ForEachBatch récepteur

Indisponible

@dp.foreach_batch_sink(...)

Non