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
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 |
| Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur. |
|
| Obligatoire. Le nom de la table ou du récepteur qui est la cible du flux d'ajout. |
|
| Le nom du flux. Si non fourni, le nom de la fonction default. |
|
| Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de
|
|
| Une description pour le flux. |
|
| Une liste de configurations Spark pour l'exécution de cette query. |
Exemples
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"))