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

Le module de pipelines n'est disponible que dans le contexte d'un pipeline. Il n'est pas disponible dans Python exécuté en dehors des pipelines. Pour plus d'informations sur la modification du code de pipeline, consultez Développer et déboguer des LakeFlow Pipelines avec l'éditeur de pipelines ETL.

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