mise à jour du flux
Utilisez le décorateur @dp.update_flow pour créer un flux de mise à jour. Les flux de mise à jour écrivent dans des puits en utilisant le mode de sortie de mise à jour, en émettant uniquement les lignes qui ont changé dans chaque batch. Contrairement aux flux d'ajout, ils prennent en charge les agrégations avec état sans nécessiter de watermark.
Les flux de mise à jour ne peuvent cibler que des récepteurs. Les tables Delta ne sont pas prises en charge.
Syntaxe
from pyspark import pipelines as dp
dp.create_sink("<sink-name>", "<format>", {"<key>": "<value>"})
@dp.update_flow(
target = "<sink-name>",
name = "<flow-name>", # optional, defaults to function name
depends_on = "<flow-name>", # optional, Public Preview
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 du récepteur auquel ce flux écrit. |
|
| Le nom du flux. Si non fourni, le nom de la fonction default. |
|
| Aperçu public. Un ou plusieurs noms de flux qui doivent s'exécuter avec succès avant le start de ce flux. Accepte un seul nom de flux ou une liste de noms. Cela ordonne uniquement l'exécution du flux ; cela ne modifie pas la façon dont le flux s'exécute. Consultez la section Ordonner l'exécution d'un flux de pipeline avec depends_on. |
|
| Une description pour le flux. |
|
| Un dictionnaire de configurations Spark pour l'exécution de cette query. Ces configurations remplacent les configurations définies pour la destination, le pipeline ou le cluster. |
Exemples
Agrégation vers un récepteur Kafka
Écrivez les résultats d'agrégation avec état dans un récepteur Kafka :
from pyspark import pipelines as dp
from pyspark.sql.functions import col
dp.create_sink("event_counts_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="event_counts_flow",
target="event_counts_sink",
)
def event_counts():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
.selectExpr("CAST(key AS STRING) AS event_type")
.groupBy(col("event_type"))
.count()
)
Mode temps réel
Aperçu public
Le mode temps réel est en Aperçu public.
Utilisez spark_conf pour configurer un flux de mise à jour pour le mode temps réel:
from pyspark import pipelines as dp
dp.create_sink("my_kafka_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="my_rtm_flow",
target="my_kafka_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
}
)
def my_real_time_flow():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
)
Limitations
- Les puits de table Delta ne sont pas pris en charge comme cibles pour les flux de mise à jour.