Aller au contenu principal

append_flow

Le décorateur @dp.append_flow crée des flux d'ajout ou des remplissages pour vos tables de pipeline. La fonction doit renvoyer un DataFrame de streaming Apache Spark. Voir Charger et traiter les données de manière incrémentielle avec LakeFlow Pipelines.

Les flux d’ajout peuvent cibler des tables de streaming, des tables gérées ou des récepteurs.

Syntaxe

Python
from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.append_flow(
target = "<target-table-name>",
name = "<flow-name>", # optional, defaults to function name
once = False, # optional
depends_on = "<flow-name>", # optional, Public Preview
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #

parameter

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

target

str

Obligatoire. Le nom de la table ou du récepteur qui est la cible du flux d'ajout.

name

str

Le nom du flux. Si non fourni, le nom de la fonction default.

once

bool

Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de once=True modifie le flux de deux manières :

  • La valeur de retour. streaming-query. Doit être un DataFrame batch dans ce cas, pas un DataFrame de streaming.
  • Le flux s'exécute une seule fois par default. Si le pipeline est mis à jour avec un refresh complet, alors le flux ONCE s'exécute de nouveau pour recréer les données.

depends_on

str OU list

Aperçu public. Un ou plusieurs noms de flux qui doivent s’exécuter avec succès avant que ce flux ne start. Accepte un nom de flux unique ou une liste de noms. Cela ordonne uniquement l’exécution du flux ; cela ne modifie pas la manière dont le flux s’exécute. Voir Ordonnancer l’exécution du flux de pipeline avec depends_on.

comment

str

Une description pour le flux.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

import_checkpoint

str

Le chemin d’accès à un point de contrôle Structured Streaming existant à importer dans le flux, afin qu’un stream migré reprenne à partir de son dernier offset validé au lieu de retraiter la source. L’importation d’un point de contrôle est en bêta. Consultez Migrate a Structured Streaming checkpoint.

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

target

str

Obligatoire. Le nom de la table ou du récepteur qui est la cible du flux d'ajout.

name

str

Le nom du flux. Si non fourni, le nom de la fonction default.

once

bool

Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de once=True modifie le flux de deux manières :

  • La valeur de retour. streaming-query. Doit être un DataFrame batch dans ce cas, pas un DataFrame de streaming.
  • Le flux s'exécute une seule fois par default. Si le pipeline est mis à jour avec un refresh complet, alors le flux ONCE s'exécute de nouveau pour recréer les données.

depends_on

str OU list

Aperçu public. Un ou plusieurs noms de flux qui doivent s’exécuter avec succès avant que ce flux ne start. Accepte un nom de flux unique ou une liste de noms. Cela ordonne uniquement l’exécution du flux ; cela ne modifie pas la manière dont le flux s’exécute. Voir Ordonnancer l’exécution du flux de pipeline avec depends_on.

comment

str

Une description pour le flux.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

import_checkpoint

str

Le chemin d’accès à un point de contrôle Structured Streaming existant à importer dans le flux, afin qu’un stream migré reprenne à partir de son dernier offset validé au lieu de retraiter la source. L’importation d’un point de contrôle est en bêta. Consultez Migrate a Structured Streaming checkpoint.

Exemples

Python
from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
return (
spark.read
.format("json")
.load("/path/to/backfill/")
)

# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

Migrer un point de contrôle Structured Streaming

info

Bêta

L’importation d’un point de contrôle est en bêta.

Utilisez import_checkpoint pour migrer une charge de travail Structured Streaming existante vers un pipeline sans retraiter la source. Définissez-le sur le checkpointLocation utilisé par votre query Structured Streaming, qui peut être un stockage cloud, un volume Unity Catalog ou un chemin DBFS. Lors de la première mise à jour du pipeline, le flux clone ce point de contrôle dans le stockage géré du pipeline. Le flux reprend ensuite à partir du dernier offset validé avec son état intact (tel que les agrégations, les clés de déduplication et les filigranes). Les mises à jour ultérieures du pipeline utilisent le point de contrôle cloné du flux ; le point de contrôle d'origine n'est pas modifié.

Le flux doit cibler une table managée créée avec create_table ou un puits de données.

Arrêtez la query Structured Streaming d’origine avant d’exécuter le pipeline. La query Structured Streaming d’origine peut être réutilisée après l’importation, mais vous devez gérer son état de point de contrôle et vous assurer que le pipeline et la query n’écrivent pas dans la même table en même temps, ce qui peut produire des données en double.

Recréez la query Structured Streaming en tant que flux de pipeline qui écrit dans une nouvelle table et importe son point de contrôle :

Python
from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
target = "target_table",
import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
# The same source your original query read from.
return spark.readStream.table("source_table")

Le point de contrôle n'est importé qu'une seule fois, lors de la première mise à jour du pipeline ; les mises à jour ultérieures ignorent import_checkpoint. Un full refresh ne réimporte pas le point de contrôle ; il start à partir d'un nouveau point de contrôle vide et retraite la source. Pour importer un autre point de contrôle, utilisez un nom de flux qui n'a pas encore été utilisé pour la table cible ; la réutilisation d'un nom de flux existant ignore l'importation.

Limitations

  • L’importation d’un point de contrôle dans une table qui existe déjà (par exemple, la query Structured Streaming d’origine) n’est pas prise en charge. Ciblez une nouvelle table créée par le pipeline ou un réceptacle.
  • import_checkpoint est pris en charge uniquement sur append_flow.