Utilisez le mode temps réel dans les LakeFlow Pipelines
Aperçu public
Le mode temps réel dans LakeFlow Pipelines est en Aperçu public sur Databricks Runtime 18.1.3. sur le canal de distribution d'aperçu.
Le mode temps réel permet un traitement des données à latence ultra-faible, avec une latence de bout en bout aussi faible que cinq millisecondes. Utilisez le mode temps réel pour les charges de travail Opérations qui nécessitent une réponse immédiate aux données de streaming, telles que la détection de fraude et la personnalisation en temps réel.
Le mode temps réel est également disponible directement dans Structured Streaming en dehors des pipelines. Voir le mode temps réel dans Structured Streaming.
Comment le mode en temps réel permet d'obtenir une faible latence
Le mode temps réel diffère du traitement continu standard de trois manières essentielles :
- Lots de longue durée : le système traite les données dès qu'elles sont disponibles dans la source au sein de lots de longue durée (la valeur par default est de cinq minutes).
- Planification simultanée des étapes : toutes les étapes de la query sont planifiées en même temps. La ressource de compute doit disposer de suffisamment d'emplacements de tâches disponibles pour couvrir toutes les étapes simultanément. Voir le dimensionnement du compute.
- Streaming shuffle : Les données sont transmises entre les étapes dès qu'elles sont produites, plutôt que d'attendre qu'une étape en amont soit terminée avant de commencer l'étape en aval.
L'intervalle de point de contrôle (configuré via pipelines.trigger.interval) contrôle la fréquence à laquelle l'état et les décalages de source sont persistés dans un stockage durable. Des intervalles plus longs réduisent la surcharge des points de contrôle, mais augmentent le temps de récupération après une défaillance et retardent le rapport des métriques. Des intervalles plus courts améliorent la durabilité, mais ajoutent des frais généraux.
Mode en temps réel et pipelines continus
Le mode temps réel est un type spécialisé de Trigger continu. Le mode continu est toujours requis ; le mode temps réel ajoute des optimisations de latence au niveau du flux en plus. Pour utiliser le mode temps réel, le pipeline doit d’abord s’exécuter en mode continu. Le mode temps réel applique ensuite des optimisations supplémentaires au niveau du flux pour atteindre une latence inférieure à la seconde, au-delà de ce qu'offre le traitement continu standard.
L'activation du mode temps réel nécessite trois étapes de configuration :
- Définissez le pipeline en mode continu.
- Activez le mode temps réel au niveau du pipeline.
- Définir un flux de mise à jour en temps réel.
Exigences
Exigence | Valeur |
|---|---|
Environnement d'exécution Databricks | 18.1.3 sur le Canal de distribution de prévisualisation des Lakeflow pipelines |
Type de compute | Classic compute ou serverless |
Configurer le mode temps réel
Étape 1 : Définissez le pipeline sur le mode continu
Dans les paramètres de votre pipeline, définissez le mode Pipeline sur Continu , ou définissez-le dans le JSON du pipeline :
{
"continuous": true
}
Étape 2 : Activer le mode temps réel au niveau du pipeline
Dans les paramètres de votre pipeline, ajoutez la clé suivante à la configuration Spark sous Avancé > Configuration Spark :
spark.databricks.streaming.realTimeMode.enabled = true
Vous pouvez également le définir dans le pipeline JSON :
{
"continuous": true,
"spark_conf": {
"spark.databricks.streaming.realTimeMode.enabled": "true"
}
}
Étape 3 : Définissez un flux de mise à jour en temps réel
Le mode en temps réel nécessite un flux de mise à jour. Utilisez dp.create_sink() pour définir la cible de sortie, puis utilisez le décorateur @dp.update_flow avec pipelines.trigger défini sur "RealTime" et target pointant vers le sink.
from pyspark import pipelines as dp
# Define the output sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "<bootstrap-servers>",
"topic": "<output-topic>",
}
)
# Define the real-time update flow targeting the sink
@dp.update_flow(
name="my_rtm_flow",
target="my_kafka_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes", # optional; defaults to 5 minutes
}
)
def my_real_time_flow():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<bootstrap-servers>")
.option("subscribe", "<input-topic>")
.load()
)
Paramètres de configuration au niveau du flux :
parameter | Obligatoire | Par défaut | Description |
|---|---|---|---|
| Oui | — | Définissez sur |
| Non |
| Intervalle de point de contrôle. Détermine la fréquence à laquelle l'état et les décalages sont validés. Des valeurs plus courtes améliorent la récupérabilité ; des valeurs plus longues réduisent les frais généraux. |
Exemples de code
Kafka à Kafka
Lisez à partir d'un sujet Kafka et écrivez vers une cible de sortie Kafka :
from pyspark import pipelines as dp
dp.create_sink("kafka_output_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="kafka_rtm_flow",
target="kafka_output_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
}
)
def kafka_rtm_flow():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.option("startingOffsets", "latest")
.load()
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "timestamp")
)
Enrichissez avec une jointure de diffusion
Joindre un Stream Kafka à une table de recherche statique. Seules les jointures de diffusion (Stream vers statique) sont prises en charge. Les jointures de Stream à Stream ne sont pas prises en charge en mode temps réel.
from pyspark import pipelines as dp
from pyspark.sql.functions import broadcast, expr
dp.create_sink("enriched_output_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": enriched_output_topic,
})
@dp.update_flow(
name="enriched_events_flow",
target="enriched_output_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
}
)
def enriched_events():
lookup = spark.read.table("catalog.schema.lookup_table")
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
.withColumn("event_key", expr("CAST(value AS STRING)"))
.join(broadcast(lookup), expr("event_key = lookup_key"))
.select("event_key", "lookup_value", "timestamp")
)
Agrégation
Comptez les événements par clé à l'aide d'un groupBy avec état. Définissez spark.sql.shuffle.partitions pour correspondre au nombre de partitions d'entrée pour les opérations avec état :
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",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
"spark.sql.shuffle.partitions": "8",
}
)
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", "timestamp")
.groupBy(col("event_type"))
.count()
)
Sources et récepteurs pris en charge
Connecteur | Comme source | En tant que puits | Notes |
|---|---|---|---|
Apache Kafka | ✓ | ✓ | — |
AWS MSK | ✓ | ✓ | Utilise l'interface compatible Kafka. |
Azure Event Hubs (connecteur Kafka) | ✓ | ✓ | Utilise l'interface compatible Kafka. |
Amazon Kinesis | ✓ | Non pris en charge | Utiliser uniquement pour le mode EFO (Enhanced Fan-Out). |
Delta | Non pris en charge | Non pris en charge | — |
Dimensionnement du compute
Vous pouvez exécuter un pipeline en temps réel par ressource de compute si le compute dispose de suffisamment d'emplacements de tâches. Les emplacements de tâches disponibles doivent couvrir toutes les tâches à travers toutes les étapes de la query.
Type de pipeline | Configuration | Emplacements de tâches requis |
|---|---|---|
Sans état à un seul étage (source Kafka + sink) |
| 8 |
À deux étapes et avec état (source Kafka + brassage) |
| 28 (8 + 20) |
À trois étapes (source Kafka + deux shuffles) |
| 48 (8 + 20 + 20) |
Si vous ne définissez pas maxPartitions, utilisez le nombre de partitions dans le sujet Kafka.
Prise en charge de l'opérateur
Catégorie | Opérateur | Pris en charge |
|---|---|---|
Sans état | Sélection, Projection | ✓ |
UDFs | UDF Scala | ✓ (avec limitations) |
UDFs | Python UDF | ✓ (avec limitations) |
Agrégation | sum, count, max, min, avg | ✓ |
Fenêtrage | Fenêtrage par basculement, Fenêtrage glissant | ✓ |
Fenêtrage | Session | Non pris en charge |
Déduplication |
| ✓ (état non délimité) |
Déduplication |
| Non pris en charge |
Jointures | Jointure de table de diffusion | ✓ |
Jointures | Jointure de Stream à Stream | Non pris en charge |
Personnalisé |
| ✓ (avec des différences comportementales) |
Personnalisé |
| ✓ (avec limitations) |
Personnalisé |
| Non pris en charge |
Personnalisé |
| Non pris en charge |
Personnalisé |
| Non pris en charge |
Personnalisé |
| Non pris en charge |
transformWithState en mode temps réel
transformWithState est pris en charge en mode temps réel avec les différences suivantes par rapport au traitement micro-batch :
handleInputRowsest invoqué une fois par ligne plutôt qu’une fois par clé et par batch. L’itérateurinputRowsrenvoie une valeur unique par appel.- Les temporisateurs d'heure d'événement ne sont pas pris en charge. Les minuteurs de temps de traitement se déclenchent lorsqu'un batch de longue durée se termine si aucune donnée n'est arrivée.
transformWithStateInPandasn'est pas pris en charge.
UDF Pandas en mode temps réel
Pour minimiser la latence avec les UDF pandas, définissez spark.sql.execution.arrow.maxRecordsPerBatch sur 1. Ceci optimise la latence au détriment du throughput. Si le throughput est également important, définissez cette valeur sur 100 ou plus.
Surveiller les performances du mode en temps réel
Le mode temps réel expose les métriques de latence dans StreamingQueryProgress sous le champ latencies. Accédez à ces métriques via un StreamingQueryListener ou en inspectant la propriété lastProgress sur la query de streaming.
Métriques | Description |
|---|---|
| Temps entre le moment où un enregistrement est lu par le flux et le moment où il est entièrement traité par le flux. |
| Temps écoulé entre le moment où un enregistrement est écrit avec succès dans le bus de messages (par exemple, le temps d'ajout du log dans Kafka) et le moment où il est lu pour la première fois par le flux |
| Latence de bout en bout totale entre la production de l'enregistrement à la source et son traitement complet par le flux |
Chaque métrique est rapportée sous forme de centiles p50, p90, p95 et p99.
Limitations
Un flux en temps réel par pipeline est recommandé. Plusieurs flux sont autorisés, mais la contention des emplacements de tâches entre les flux augmente la latence.
Pour une liste complète des limitations des opérateurs et des sources, consultez Limitations du mode temps réel.