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
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 |
| 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. |
|
| 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. |
|
| Un chemin de point de contrôle externe à importer avant de démarrer le flux. Importé une seule fois, lorsque le répertoire de checkpoint du flux n'existe pas encore. |
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.