Didacticiel : Créer plusieurs flux avec des paramètres différents
Un pipeline peut contenir plusieurs flux qui sont presque identiques, ne différant que par quelques paramètres. Définir ces flux explicitement est sujet aux erreurs, redondant et difficile à maintenir. La métaprogrammation avec les fonctions internes Python génère des flux répétitifs dynamiquement, chaque invocation fournissant un ensemble différent de paramètres.
Présentation de la métaprogrammation
La métaprogrammation dans les LakeFlow Pipelines utilise les fonctions internes de Python. Étant donné que ces fonctions sont évaluées de manière paresseuse par le runtime du pipeline, vous pouvez envelopper les décorateurs @dp.table dans une fonction fabrique et appeler cette fabrique plusieurs fois avec des arguments différents. Chaque appel enregistre un nouveau flux sans dupliquer le code.
Pour plus de détails sur l'utilisation des boucles for avec les LakeFlow Pipelines, consultez Créer des tables dans une boucle for.
Exemple : délais d'intervention des pompiers
L'exemple suivant utilise le dataset des services d'incendie intégré afin de trouver les quartiers avec les temps de réponse d'urgence les plus rapides pour chaque type d'appel. Sans métaprogrammation, vous devez écrire des définitions de table presque identiques pour chaque type d'appel (Alarmes, Incendie de structure, Incident médical). Avec la métaprogrammation, une seule fonction d'usine les génère toutes.
Étape 1 : définissez la table d'ingestion brute
import functools
from pyspark import pipelines as dp
from pyspark.sql.functions import *
@dp.table(
name="raw_fire_department",
comment="raw table for fire department response"
)
@dp.expect_or_drop("valid_received", "received IS NOT NULL")
@dp.expect_or_drop("valid_response", "responded IS NOT NULL")
@dp.expect_or_drop("valid_neighborhood", "neighborhood != 'None'")
def get_raw_fire_department():
return (
spark.read.format('csv')
.option('header', 'true')
.option('multiline', 'true')
.load('/databricks-datasets/timeseries/Fires/Fire_Department_Calls_for_Service.csv')
.withColumnRenamed('Call Type', 'call_type')
.withColumnRenamed('Received DtTm', 'received')
.withColumnRenamed('Response DtTm', 'responded')
.withColumnRenamed('Neighborhooods - Analysis Boundaries', 'neighborhood')
.select('call_type', 'received', 'responded', 'neighborhood')
)
Étape 2 : Définir la fonction de fabrique de flux
La fonction de fabrique generate_tables enregistre deux tables pour chaque type d'appel : une table d'appels filtrés et une table de temps de réponse classés. Les deux sont créées comme fonctions internes décorées avec @dp.table.
all_tables = []
def generate_tables(call_table, response_table, filter):
@dp.table(
name=call_table,
comment="top level tables by call type"
)
def create_call_table():
return spark.sql("""
SELECT
unix_timestamp(received,'M/d/yyyy h:m:s a') as ts_received,
unix_timestamp(responded,'M/d/yyyy h:m:s a') as ts_responded,
neighborhood
FROM raw_fire_department
WHERE call_type = '{filter}'
""".format(filter=filter))
@dp.table(
name=response_table,
comment="top 10 neighborhoods with fastest response time"
)
def create_response_table():
return spark.sql("""
SELECT
neighborhood,
AVG((ts_received - ts_responded)) as response_time
FROM {call_table}
GROUP BY 1
ORDER BY response_time
LIMIT 10
""".format(call_table=call_table))
all_tables.append(response_table)
Étape 3 : Invoquer la fabrique et définir la table récapitulative
Appelez l'usine une fois pour chaque type d'appel, puis définissez un tableau récapitulatif qui unit les résultats pour trouver les voisinages qui apparaissent le plus souvent dans toutes les catégories.
generate_tables("alarms_table", "alarms_response", "Alarms")
generate_tables("fire_table", "fire_response", "Structure Fire")
generate_tables("medical_table", "medical_response", "Medical Incident")
@dp.table(
name="best_neighborhoods",
comment="which neighbor appears in the best response time list the most"
)
def summary():
target_tables = [dp.read(t) for t in all_tables]
unioned = functools.reduce(lambda x, y: x.union(y), target_tables)
return (
unioned.groupBy(col("neighborhood"))
.agg(count("*").alias("score"))
.orderBy(desc("score"))
)
Après avoir exécuté ce pipeline, vous créez un ensemble de tables similaires, comme ce graphe :

Concepts clés
- Les fonctions internes sont enregistrées paresseusement : le décorateur
@dp.tablen'exécute pas la fonction immédiatement. Il enregistre la fonction avec le runtime du pipeline, ce qui résout le Graphe complet du flux de données avant que l'exécution ne commence. - Les fermetures capturent les paramètres : Chaque fonction interne s'applique aux paramètres transmis à la fabrique (
call_table,response_table,filter), de sorte que chaque flux enregistré utilise son propre ensemble de valeurs isolé. - **Listes de tables dynamiques** : L'utilisation d'une liste comme
all_tablespour suivre les noms de tables générés par programme facilite leur référencement ultérieur (par exemple, dans une union ou une jointure).