Aller au contenu principal

Utiliser les flux dans les LakeFlow Pipelines

Les flux dans un LakeFlow Pipelines déplacent les données vers une table de streaming ou une vue matérialisée. Les exemples suivants montrent comment définir des flux default, définir un flux séparément de sa cible, écrire dans une table de streaming à partir de plusieurs rubriques Kafka, exécuter un remplissage rétroactif unique, et remplacer les queries UNION par un traitement de flux d'ajout.

Pour une vue d'ensemble des flux, consultez Charger et traiter les données de manière incrémentielle avec LakeFlow Pipelines.

Exemple : Créez un flow par default

Lorsque vous créez un pipeline, vous définissez généralement une table ou une vue ainsi que la query qui la prend en charge. Par exemple, cette query crée une table de streaming nommée customers_silver en lisant depuis customers_bronze. La table de streaming et son flux default sont créés ensemble en une seule étape.

SQL
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)

Le flux default pour une table de streaming est un flux d' ajout qui ajoute de nouvelles lignes à chaque mise à jour, et il porte le même nom que la cible. C'est la manière la plus courante d'utiliser les pipelines — en créant un flux et sa cible en une seule étape — et vous pouvez l'utiliser pour ingérer ou transformer des données. Pour en savoir plus sur les concepts de flux, consultez Charger et traiter les données de manière incrémentielle avec les LakeFlow Pipelines.

Exemple : Définir un flux séparément de sa cible

Vous pouvez également créer un flux pour une table que vous avez définie séparément. Le résultat est identique à la création d'un flux default, y compris l'utilisation du même nom pour la table de streaming et le flux :

Python
from pyspark import pipelines as dp

# create streaming table
dp.create_streaming_table("customers_silver")

# add a flow
@dp.append_flow(
target = "customers_silver")
def customer_silver():
return spark.readStream.table("customers_bronze")

Définir un flux séparément de sa cible vous permet de créer plusieurs flux qui ajoutent des données à la même cible. Utilisez le décorateur @dp.append_flow dans l'interface Python ou la clause CREATE FLOW...INSERT INTO dans l'interface SQL pour ajouter des flux pour des tâches telles que les suivantes :

Pour les query Python, utilisez la fonction create_streaming_table() pour créer une table cible.

important
  • Si vous devez définir des contraintes de qualité des données avec des attentes, définissez les attentes sur la table cible dans le cadre de la fonction create_streaming_table() ou sur une définition de table existante. Vous ne pouvez pas définir d'attentes dans la définition @append_flow.
  • Les flux sont identifiés par un nom de flux , et ce nom est utilisé pour identifier les points de contrôle en streaming. L'utilisation du nom de flux pour identifier le point de contrôle signifie ce qui suit :
    • Si un flux existant dans un pipeline est renommé, le point de contrôle n'est pas reporté, et le flux renommé est effectivement un flux entièrement nouveau.
    • Vous ne pouvez pas réutiliser un nom de flux dans un pipeline, car le point de contrôle existant ne correspondra pas à la nouvelle définition de flux.

Exemple : Écrire dans une table de streaming à partir de plusieurs rubriques Kafka

Les exemples suivants créent une table de streaming nommée kafka_target et écrivent dans cette table de streaming à partir de deux sujets Kafka :

Python
from pyspark import pipelines as dp

dp.create_streaming_table("kafka_target")

# Kafka stream from multiple topics
@dp.append_flow(target = "kafka_target")
def topic1():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic1")
.load()
)

@dp.append_flow(target = "kafka_target")
def topic2():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic2")
.load()
)

Pour en savoir plus sur la fonction table à valeur read_kafka() utilisée dans les requêtes SQL, consultez read_kafka dans la référence du langage SQL.

En Python, vous pouvez créer par programme plusieurs flux qui ciblent une seule table. L'exemple suivant montre ce modèle pour une liste de sujets Kafka.

remarque

Ce modèle a les mêmes exigences que l'utilisation d'une boucle for pour créer des tables. Vous devez explicitement transmettre une valeur Python à la fonction définissant le flux. Consultez Créer des tables dans une boucle for.

Python
from pyspark import pipelines as dp

dp.create_streaming_table("kafka_target")

topic_list = ["topic1", "topic2", "topic3"]

for topic_name in topic_list:

@dp.append_flow(target = "kafka_target", name=f"{topic_name}_flow")
def topic_flow(topic=topic_name):
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", topic)
.load()
)

Exemple : exécuter un remplissage de données ponctuel

Si vous voulez exécuter une query pour ajouter des données à une table de streaming existante, utilisez append_flow.

Après avoir ajouté un ensemble de données existantes, vous disposez de plusieurs options :

  • Si vous souhaitez que la query ajoute de nouvelles données si elles arrivent dans le répertoire de rattrapage, laissez la query en place.
  • Si vous souhaitez que ce soit un backfill unique et qu'il ne s'exécute plus jamais, supprimez la query après avoir exécuté le pipeline une seule fois.
  • Si vous souhaitez que la requête s'exécute une seule fois, et ne s'exécute à nouveau que si les données sont entièrement actualisées, définissez le paramètre once sur True sur le flux d'ajout. En SQL, utilisez INSERT INTO ONCE.

Les exemples suivants exécutent une query pour ajouter des données historiques à une table de streaming :

Python
from pyspark import pipelines as dp

@dp.table()
def csv_target():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/sourceDir")

@dp.append_flow(
target = "csv_target",
once = True)
def backfill():
return spark.read
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/backfill/data/dir")

Pour un exemple plus détaillé, consultez Remplir les données historiques avec des pipelines.

Exemple : Utiliser le traitement de flux d'ajout au lieu de UNION

Au lieu d'utiliser une query avec une clause UNION, vous pouvez utiliser des requêtes de flux d'ajout pour combiner plusieurs sources et écrire dans une seule table de streaming. L'utilisation de requêtes de flux d'ajout, au lieu de UNION, vous permet d'ajouter à une table de streaming à partir de plusieurs sources sans exécuter un refresh complet.

L'exemple Python suivant inclut une query qui combine plusieurs sources de données avec une clause UNION :

Python
@dp.create_table(name="raw_orders")
def unioned_raw_orders():
raw_orders_us = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")
)

raw_orders_eu = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")
)

return raw_orders_us.union(raw_orders_eu)

Les exemples suivants remplacent la query UNION par des query de flux d'ajout :

Python
dp.create_streaming_table("raw_orders")

@dp.append_flow(target="raw_orders")
def raw_orders_us():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")

@dp.append_flow(target="raw_orders")
def raw_orders_eu():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")

# Additional flows can be added without the full refresh that a UNION query would require:
@dp.append_flow(target="raw_orders")
def raw_orders_apac():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/apac")

Exemple : utilisez transformWithState pour surveiller les pulsations des capteurs

L'exemple suivant présente un processeur avec état qui lit depuis Kafka et vérifie que les capteurs émettent des signaux de présence périodiquement. Si un signal de présence n'est pas reçu dans les 5 minutes, le processeur émet une entrée vers la table Delta cible pour l'analyse.

Pour plus d'informations sur la création d'applications avec état personnalisées, consultez Créez une application avec état personnalisée.

remarque

RocksDB est le fournisseur d'état par default à partir de Databricks Runtime 17.2. Si la query échoue en raison d'une exception de fournisseur non prise en charge, ajoutez les configurations de pipeline suivantes, effectuez un refresh complet ou une Reset du checkpoint, puis réexécutez votre pipeline :

JSON
"configuration": {
"spark.sql.streaming.stateStore.providerClass": "com.databricks.sql.streaming.state.RocksDBStateStoreProvider",
"spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled": "true"
}
Python
from typing import Iterator

import pandas as pd

from pyspark import pipelines as dp
from pyspark.sql.functions import col, from_json
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType, StringType, TimestampType

KAFKA_TOPIC = "<your-kafka-topic>"

output_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])

class SensorHeartbeatProcessor(StatefulProcessor):
def init(self, handle: StatefulProcessorHandle) -> None:
# Define state schema to store sensor information (sensor_id is the grouping key)
state_schema = StructType([
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])
self.sensor_state = handle.getValueState("sensorState", state_schema)
# State variable to track the previously registered timer
timer_schema = StructType([StructField("timer_ts", LongType(), False)])
self.timer_state = handle.getValueState("timerState", timer_schema)
self.handle = handle

def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
# Process one row from input and update state
pdf = next(rows)
row = pdf.iloc[0]
# Store or update the sensor information in state using current timestamp
current_time = pd.Timestamp(timerValues.getCurrentProcessingTimeInMs(), unit='ms')
self.sensor_state.update((
row["sensor_type"],
current_time
))

# Delete old timer if already registered
if self.timer_state.exists():
old_timer = self.timer_state.get()[0]
self.handle.deleteTimer(old_timer)

# Register a timer for 5 minutes from current processing time
expiry_time = timerValues.getCurrentProcessingTimeInMs() + (5 * 60 * 1000)
self.handle.registerTimer(expiry_time)
# Store the new timer timestamp in state
self.timer_state.update((expiry_time,))

# No output on input processing, output only on timer expiry
return iter([])

def handleExpiredTimer(self, key, timerValues, expiredTimerInfo) -> Iterator[pd.DataFrame]:
# Emit output row based on state store
if self.sensor_state.exists():
state = self.sensor_state.get()
output = pd.DataFrame({
"sensor_id": [key[0]], # Use grouping key as sensor_id
"sensor_type": [state[0]],
"last_heartbeat_time": [state[1]]
})
# Remove the entry for the sensor from the state store
self.sensor_state.clear()
# Remove the timer state entry
self.timer_state.clear()
yield output

def close(self) -> None:
pass

dp.create_streaming_table("sensorAlerts")

# Define the schema for the Kafka message value
sensor_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("sensor_value", LongType(), False)])

@dp.append_flow(target = "sensorAlerts")
def kafka_delta_flow():
return (
spark.readStream
.format("kafka")
.option("subscribe", KAFKA_TOPIC)
.option("startingOffsets", "earliest")
.load()
.select(from_json(col("value").cast("string"), sensor_schema).alias("data"), col("timestamp"))
.select("data.*", "timestamp")
.withWatermark('timestamp', '1 hour')
.groupBy(col("sensor_id"))
.transformWithStateInPandas(
statefulProcessor = SensorHeartbeatProcessor(),
outputStructType = output_schema,
outputMode = 'update',
timeMode = 'ProcessingTime'))