Aller au contenu principal

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

Python
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

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 du récepteur auquel ce flux écrit.

name

str

Le nom du flux. Si non fourni, le nom de la fonction default.

comment

str

Une description pour le flux.

spark_conf

dict

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.

import_checkpoint

str

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.

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 du récepteur auquel ce flux écrit.

name

str

Le nom du flux. Si non fourni, le nom de la fonction default.

comment

str

Une description pour le flux.

spark_conf

dict

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.

import_checkpoint

str

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 :

Python
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

info

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:

Python
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={
&quot;pipelines.trigger&quot;: &quot;RealTime&quot;,
&quot;pipelines.trigger.interval&quot;: &quot;5 minutes&quot;,
}
)
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.