Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Le décorateur @dp.append_flow crée des flux d'ajout ou des backfills pour vos tables de pipeline. La fonction doit retourner un DataFrame de streaming Apache Spark. Consultez Charger et traiter des données de manière incrémentielle avec des flux de pipeline Lakeflow.
Les flux d’ajout peuvent cibler des tables ou récepteurs de streaming.
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
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #
Paramètres
| Paramètre | Type | Descriptif |
|---|---|---|
| function | function |
Obligatoire. Fonction qui retourne un DataFrame de streaming Apache Spark à partir d’une requête définie par l’utilisateur. |
target |
str |
Obligatoire. Nom de la table ou du récepteur qui est la cible du flux d’ajout. |
name |
str |
Nom du flux. S’il n’est pas fourni, la valeur par défaut est le nom de la fonction. |
once |
bool |
Si vous le souhaitez, définissez le flux en tant que flux à usage unique, tel qu’un remblai. L'utilisation de once=True change le flux de deux manières :
|
comment |
str |
Description du flux. |
spark_conf |
dict |
Liste des configurations Spark pour l’exécution de cette requête |
import_checkpoint |
str |
Le chemin vers un point de contrôle de Structured Streaming existant pour importer dans le flux, afin qu’un flux migré reprenne à partir de son dernier décalage engagé au lieu de retraiter la source. L’importation d’un point de contrôle est en bêta. Voir Migrer un point de contrôle de streaming structuré. |
Examples
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 de streaming structuré
Important
L’importation d’un point de contrôle est en bêta.
Utilisez import_checkpoint pour migrer une charge de travail de streaming structuré existante vers un pipeline sans retraiter la source. Régle-la sur la checkpointLocation requête de Structured Streaming utilisée, qui peut être un stockage cloud, un volume Unity Catalog ou un chemin DBFS. Lors de la première mise à jour du pipeline, le flow clone ce point de contrôle dans le stockage géré du pipeline. Le flux reprend alors à partir du dernier décalage engagé avec son état (comme les agrégations, les clés de déduplication et les filigranes) intact. Les mises à jour ultérieures du pipeline utilisent le point de contrôle cloné du flux ; le point de contrôle original n’est pas modifié.
Le flux doit viser une table managée créée avec create_table ou un puits d’eau.
Arrêtez la requête originale de Structured Streaming avant de lancer le pipeline. La requête Structured Streaming originale peut être réutilisée après l’importation, mais il faut gérer son état du point de contrôle et s’assurer que le pipeline et la requête n’écrivent pas simultanément sur la même table, ce qui peut produire des données dupliquées.
Recréez la requête de Structured Streaming comme un flux pipeline qui écrit sur une nouvelle table et importe son point de contrôle :
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 rafraîchissement complet ne réimporte pas le point de contrôle ; il commence depuis 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 jamais été utilisé pour la table cible auparavant ; réutiliser un nom de flux existant saute l’importation.
Limitations
- L’importation d’un point de contrôle dans une table déjà existante (par exemple, la cible originale de requête Structured Streaming) n’est pas prise en charge. Ciblez une nouvelle table créée par le pipeline, ou un puits d’envoi.
-
import_checkpointest uniquement supporté sur append_flow.