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
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

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.

depends_on

str OU list

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.

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.

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.

depends_on

str OU list

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.

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.

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.