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 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
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>") # 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.

comment

str

Une description pour le flux.

spark_conf

dict

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

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.

comment

str

Une description pour le flux.

spark_conf

dict

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

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"))