Architecture en étoile et en éventail dans les LakeFlow Pipelines
Le fan-in et le fan-out sont des modèles courants pour la création de pipeline de données évolutives et fiables. Cette page explique les deux et montre comment les implémenter dans les Lakeflow pipelines.
Que sont le fan-in et le fan-out ?
**Le Fan-in** est un modèle architectural où les données de multiples sources sont ingérées et traitées au sein d'un seul pipeline.

Les sources peuvent inclure :
- Stream d'événements en temps réel (par exemple, Kafka et Kinesis)
- Stockage cloud (par exemple, S3, ADLS, et Google Cloud Storage)
- Bases de données relationnelles (par exemple, PostgreSQL, MySQL et Snowflake)
- appareils IoT (par exemple, capteurs, logs et API)
En consolidant divers flux de Stream dans une seule couche de traitement, le fan-in permet une Transformations cohérente, une déduplication et un enrichissement de données avant que les données ne se déplacent en aval.
**La diffusion en éventail** suit une approche un à plusieurs, acheminant un unique et traité de données Stream vers plusieurs destinations.

Les destinations peuvent inclure :
- Tables Delta pour le stockage structuré
- Systèmes d'alerte en temps réel pour la détection d'anomalies
- Modèles de machine learning pour l'analytique prédictive
- Data warehouses pour le reporting et l'analytique
- Files d'attente de messages pour la communication asynchrone et le traitement découplé.
Ce modèle garantit que chaque système en aval reçoit les données dans le format requis, ce qui permet aux organisations d'intégrer les données de streaming dans diverses applications métier.
En pratique, les pipelines combinent souvent les deux modèles. Par exemple :
- Une entreprise collecte les données d'activité des utilisateurs à partir de multiples applications, sites web et appareils mobiles (fan-in).
- Les données traitées sont stockées dans Delta Lake pour l'analyse historique tandis que des alertes en temps réel Trigger pour une activité inhabituelle (fan-out).
Implémenter le fan-in avec des flux d'ajout
Les pipelines en entonnoir Mergent plusieurs données Stream en une cible unifiée. Traditionnellement, cela nécessite des queries d'union complexes et la création manuelle de points de contrôle. Les flux d'ajout simplifient cela en permettant à divers flux de données d'alimenter directement une seule table de streaming sans unions explicites ni logique complexe. Chaque source est gérée indépendamment, permettant l'ingestion incrémentielle des données et les mises à jour.
Par exemple, utilisez des flux d'ajout pour consolider plusieurs rubriques Kafka ou des Stream de données régionaux dans une table cible unifiée.
- Python
- SQL
from pyspark import pipelines as dp
dp.create_streaming_table("all_topics")
# Kafka stream from topic1
@dp.append_flow(target="all_topics")
def topic1():
return spark.readStream.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,...") \
.option("subscribe", "topic1") \
.load()
# Kafka stream from topic2
@dp.append_flow(target="all_topics")
def topic2():
return spark.readStream.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,...") \
.option("subscribe", "topic2") \
.load()
CREATE OR REFRESH STREAMING TABLE all_topics;
CREATE FLOW
topic1
AS INSERT INTO
all_topics BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic1');
CREATE FLOW
topic2
AS INSERT INTO
all_topics BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic2');
Mettre en œuvre le fan-out
Les pipelines de distribution en éventail distribuent les données d'une source vers plusieurs sorties. LakeFlow Pipelines supportent trois approches selon votre cas d'utilisation.
Utilisez des boucles for pour une logique généralisée
Si votre logique ETL est identique pour plusieurs cibles, utilisez des boucles Python pour générer dynamiquement plusieurs tables par le biais de boucles paramétrées. Ceci évite le codage répétitif et simplifie la mise à l'échelle du pipeline par la configuration.
Chaque flux ou tableau généré traite l'intégralité du dataset source de manière indépendante. Pour les sources avec des limites de throughput partagé ou de capacité de lecture, telles que Kafka, cela peut avoir un impact significatif sur les performances. Évaluez attentivement l'approche pour de telles sources avant de l'utiliser.
regions = ["US", "EU", "APAC"]
for region in regions:
@dp.materialized_view(name=f"orders_{region.lower()}_filtered")
def filtered_orders(region_filter=region):
return spark.read.table("combined_orders").filter(f"region = '{region_filter}'")
Utiliser des flux indépendants pour une logique spécifique à la cible
Lorsque les transformations ETL varient considérablement par cible, mettez en œuvre des flux de données indépendants. Cette approche offre un contrôle précis et des performances optimisées, adaptées à chaque cas d'utilisation.
from pyspark import pipelines as dp
# Grouped output
@dp.materialized_view(name="orders_sink")
def region_orders():
df = spark.read.table("combined_orders").groupBy("region").count()
# Add additional logic here
return df
# BI materialized view
@dp.materialized_view(name="orders_bi_materialized")
def orders_bi():
return spark.read.table("combined_orders").select("order_id", "amount", "region")
# ML feature table
@dp.materialized_view(name="orders_ml_features")
def orders_ml():
return (
spark.read.table("combined_orders")
.withColumn("high_value_order", col("amount") > 1000)
.select("order_id", "high_value_order", "region")
)
Utilisez ForEachBatch pour le routage personnalisé.
Aperçu public
foreach_batch_sink est disponible en aperçu public via le Canal de distribution de l'aperçu LakeFlow Pipelines. Consultez channel dans les configurations de pipeline.
Le foreach_batch_sink applique une logique personnalisée à chaque micro-lot, permettant des Transformations complexes, la fusion ou le routage vers plusieurs destinations, y compris celles sans prise en charge intégrée du streaming, telles que les récepteurs JDBC.
Chaque batch exécute plusieurs opérations d'écriture indépendamment. Les échecs d'une opération n'annulent pas automatiquement les écritures précédemment réussies. Cela peut entraîner des données partielles ou incohérentes entre les cibles, en particulier lors du traitement de sources partagées comme Kafka. Concevez vos pipelines avec une gestion minutieuse des erreurs et des tests approfondis. Voir Utiliser ForEachBatch pour écrire dans des puits de données arbitraires dans les pipelines.
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="user_events_feb")
def user_events_handler(batch_df, batch_id):
# Write to Delta table
batch_df.write.format("delta").mode("append").saveAsTable("my_catalog.my_schema.my_delta_table")
# Write to JSON files
batch_df.write.format("json").mode("append").save("/Volumes/path/to/json_target")
@dp.append_flow(target="user_events_feb", name="user_events_flow")
def read_user_events():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/data/incoming/events")
)
Modèles ForEachBatch courants
Le foreach_batch_sink prend en charge plusieurs modèles. Certains schémas courants incluent :
-
Flux unique vers récepteur multidestination : Un seul
append_flowlit à partir d'une source de streaming et achemine les données vers unforeach_batch_sink. Le récepteur gère l'écriture vers plusieurs destinations (par exemple, Delta, JSON et systèmes externes). Ceci est idéal pour les cas d'utilisation à sorties multiples simples avec une logique de transformations partagée. -
**Plusieurs flux vers un récepteur unifié** : Plusieurs
append_flowsources (par exemple, différents répertoires, formats, sujets Kafka ou APIs externes) fusionnent en unforeach_batch_sinkseul. Ceci centralise la logique de transformation courante, la gestion des sorties et la gestion des erreurs. Comme un seul point de contrôle doit être maintenu, cette approche réduit considérablement la complexité de coordination. C'est particulièrement utile lors du traitement des files de messages tels que Kafka ou des APIs externes. -
Un flux pour un récepteur (plusieurs paires indépendantes) : Chaque
append_flowa unforeach_batch_sinkdédié, établissant des relations claires et isolées entre les sources uniques et leurs cibles. Ceci est idéal pour les pipelines avec de nombreux flux indépendants nécessitant une logique de traitement unique, un dépannage simplifié et une gestion des erreurs isolée.
En pratique, ces approches se complètent souvent. Par exemple, utilisez des boucles pour générer dynamiquement plusieurs flux d’ajout pour les scénarios d’agrégation à grande échelle, puis distribuez les résultats à l’aide de boucles ou de foreach_batch_sink pour la diffusion.
Bonnes pratiques
- Les flux d'ajout requièrent que les schémas source s'alignent avec la table de streaming cible pour éviter les erreurs de traitement. Utilisez les attentes de schéma de LakeFlow Pipelines pour détecter et gérer les exceptions de manière proactive, garantissant la cohérence du schéma tout au long du pipeline.
- Gardez la logique de boucle for bien définie et simple.
- Nommez chaque flux et chaque table clairement pour maintenir la lisibilité.
- Surveiller l'utilisation des ressources pour monter en charge efficacement et éviter les goulots d'étranglement de performance.
- Lorsque vous écrivez dans des files d'attente de messages, utilisez un
foreach_batch_sinkavec un seulappend_flowqui consolide tous les Stream d'entrée. Cela simplifie la gestion des états en aval et des points de contrôle.
Limitations
- L'interface utilisateur de lignage des LakeFlow Pipelines peut ne pas afficher les métriques et les métadonnées au niveau du flux pour les nouvelles sources de flux d'ajout.
- Développez plutôt que de réduire la liste des valeurs utilisées dans une boucle for. Si un dataset précédemment défini est omis lors des exécutions ultérieures du pipeline, il est automatiquement supprimé du schéma cible, ce qui entraîne une perte de données involontaire.