Transformer des données avec des pipelines
Vous déclarez des Transformations dans les pipelines pour spécifier la manière dont les enregistrements sont traités par la logique de query, en utilisant des modèles courants tels que les jointures stream-statique, les agrégations incémentielles et le mélange de tables de streaming et de vues matérialisées.
Vous pouvez définir un dataset par rapport à toute query qui renvoie un DataFrame. Vous pouvez utiliser les opérations intégrées d'Apache Spark, les UDF, la logique personnalisée et les modèles MLflow comme transformations dans votre pipeline. Une fois les données ingérées dans votre pipeline, vous pouvez définir de nouveaux datasets par rapport aux sources en amont pour créer de nouvelles tables de streaming, vues matérialisées et vues.
Pour obtenir des conseils sur le choix entre les vues, les vues matérialisées et les tables de streaming, consultez Qu'est-ce qu'un pipeline ?. Pour savoir comment effectuer efficacement un traitement avec état dans un pipeline, consultez Optimiser le traitement avec état à l'aide de filigranes.
Exclure les tables du schéma cible
Si vous devez calculer des tables intermédiaires non destinées à la consommation externe, vous pouvez empêcher leur publication dans un schéma à l'aide du mot-clé PRIVATE. Les tables privées stockent et traitent toujours les données conformément à la sémantique du pipeline, mais ne doivent pas être accédées en dehors du pipeline actuel. Une table privée persiste pendant toute la durée de vie du pipeline qui la crée. Utilisez la syntaxe suivante pour déclarer des tables privées :
- SQL
- Python
CREATE PRIVATE STREAMING TABLE private_table
AS SELECT ... ;
@dp.table(
private=True)
def private_table():
return ("...")
Combiner des tables en streaming et des vues matérialisées dans un seul pipeline
Les tables de streaming héritent des garanties de traitement d’Apache Spark Structured Streaming et sont configurées pour traiter les query provenant de sources de données en mode ajout uniquement, où les nouvelles lignes sont toujours insérées dans la table source plutôt que modifiées.
Although, by default, streaming tables require append-only data sources, when a streaming source is another streaming table that requires updates or deletes, you can override this behavior with the skipChangeCommits flag
Un schéma de streaming courant implique l'ingestion de données source pour créer les datasets initiaux dans un pipeline. Ces datasets initiaux sont communément appelés tables bronze et effectuent souvent des transformations simples.
En revanche, les tables finales d'un pipeline, communément appelées tables Gold, nécessitent souvent des agrégations complexes ou une lecture à partir de cibles d'une opération AUTO CDC ... INTO. Étant donné que ces opérations créent intrinsèquement des mises à jour plutôt que des ajouts, elles ne sont pas prises en charge comme entrées pour les tables de streaming. Ces transformations sont mieux adaptées aux vues matérialisées.
En mélangeant des tables de streaming et des vues matérialisées dans un seul pipeline, vous pouvez simplifier votre pipeline, éviter une nouvelle ingestion ou un nouveau traitement coûteux des données brutes, et disposer de toute la puissance de SQL pour calculer des agrégations complexes sur un dataset encodé et filtré de manière efficace. L'exemple suivant illustre ce type de traitement mixte :
Ces exemples utilisent Auto Loader pour charger des fichiers à partir du stockage cloud. Pour charger des fichiers avec Auto Loader dans un pipeline Unity Catalog activé, vous devez utiliser des emplacements externes. Pour en savoir plus sur l'utilisation d'Unity Catalog avec des pipelines, consultez Utiliser Unity Catalog avec des pipelines.
- Python
- SQL
@dp.table
def streaming_bronze():
return (
# Since this is a streaming source, this table is incremental.
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("s3://path/to/raw/data")
)
@dp.table
def streaming_silver():
# Since we read the bronze table as a stream, this silver table is also
# updated incrementally.
return spark.readStream.table("streaming_bronze").where(...)
@dp.materialized_view
def live_gold():
# This table will be recomputed completely by reading the whole silver table
# when it is updated.
return spark.read.table("streaming_silver").groupBy("user_id").count()
CREATE OR REFRESH STREAMING TABLE streaming_bronze
AS SELECT * FROM STREAM read_files(
"s3://path/to/raw/data",
format => "json"
)
CREATE OR REFRESH STREAMING TABLE streaming_silver
AS SELECT * FROM STREAM(streaming_bronze) WHERE...
CREATE OR REFRESH MATERIALIZED VIEW mv_gold
AS SELECT count(*) FROM streaming_silver GROUP BY user_id
Découvrez-en davantage sur l’utilisation d’ Auto Loader pour ingérer de façon incrémentielle des fichiers JSON depuis S3.
Jointures Stream-statiques
Les jointures Stream-static sont un bon choix pour dénormaliser un flux continu de données en ajout seulement avec une table de dimension principalement statique.
À chaque mise à jour du pipeline, les nouveaux enregistrements du Stream sont joints à l'instantané le plus récent de la table statique. Si des enregistrements sont ajoutés ou mis à jour dans la table statique après le traitement des données correspondantes de la table de streaming, les enregistrements résultants ne sont pas recalculés, sauf si un refresh complet est effectué.
Dans les pipelines configurés pour une exécution Trigger, la table statique renvoie les résultats à partir du moment où la mise à jour a start. Dans les pipelines configurés pour une exécution continue, la version la plus récente de la table statique est interrogée chaque fois que la table traite une mise à jour.
Voici un exemple de jointure stream-statique :
- Python
- SQL
@dp.table
def customer_sales():
return spark.readStream.table("sales").join(spark.read.table("customers"), ["customer_id"], "left")
CREATE OR REFRESH STREAMING TABLE customer_sales
AS SELECT * FROM STREAM(sales)
INNER JOIN LEFT customers USING (customer_id)
Calculer les agrégats efficacement
Vous pouvez utiliser des tables de streaming pour calculer de manière incrémentielle des agrégats distributifs simples comme le nombre, le min, le max ou la somme, et des agrégats algébriques comme la moyenne ou l'écart type. Databricks recommande l'agrégation incrémentielle pour les queries avec un nombre limité de groupes, telles qu'une query avec une clause GROUP BY country. Seules les nouvelles données d'entrée sont lues à chaque mise à jour.
Pour en savoir plus sur l'écriture de requêtes de pipeline qui effectuent des agrégations incrémentielles, consultez Effectuer des agrégations fenêtrées avec des filigranes.
Utiliser les modèles MLflow dans les pipelines
Pour utiliser des modèles MLflow dans un pipeline compatible Unity Catalog, votre pipeline doit être configuré pour utiliser le canal preview. Pour utiliser le current Canal de distribution, vous devez configurer votre pipeline pour publier dans le Hive metastore.
Vous pouvez utiliser les modèles entraînés avec MLflow dans les pipelines. Les modèles MLflow sont traités comme des transformations dans Databricks, ce qui signifie qu'ils agissent sur une entrée de DataFrame Spark et renvoient les résultats sous forme de DataFrame Spark. Puisque les pipelines définissent des datasets par rapport aux DataFrames, vous pouvez convertir les charges de travail Apache Spark qui utilisent MLflow en pipelines avec seulement quelques lignes de code. Pour en savoir plus sur MLflow, consultez MLflow sur Databricks.
Si vous disposez déjà d'un script Python appelant un modèle MLflow, vous pouvez adapter ce code à un pipeline en utilisant le décorateur @dp.table ou @dp.materialized_view et en vous assurant que les fonctions sont définies pour renvoyer les résultats des transformations. Les pipelines n'installent pas MLflow par default, veuillez donc confirmer que vous avez installé les bibliothèques MLFlow avec %pip install mlflow et importé mlflow et dp en haut de votre source. Pour une introduction à la syntaxe du pipeline, consultez Développer du code de pipeline avec Python.
Pour utiliser les modèles MLflow dans les pipelines, suivez les étapes suivantes :
- Obtenez l’ID d’exécution et le nom du modèle MLflow. L’ID d’exécution et le nom du modèle sont utilisés pour construire l’URI du modèle MLflow.
- Utilisez l'URI pour définir une UDF Spark afin de charger le modèle MLflow.
- Appelez l'UDF dans vos définitions de table pour utiliser le modèle MLflow.
L'exemple suivant montre la syntaxe de base pour ce modèle :
%pip install mlflow==2.20.2
from pyspark import pipelines as dp
import mlflow
run_id= "<mlflow-run-id>"
model_name = "<the-model-name-in-run>"
model_uri = f"runs:/{run_id}/{model_name}"
loaded_model_udf = mlflow.pyfunc.spark_udf(spark, model_uri=model_uri)
@dp.materialized_view
def model_predictions():
return spark.read.table(<input-data>)
.withColumn("prediction", loaded_model_udf(<model-features>))
À titre d'exemple complet, le code suivant définit une UDF Spark nommée loaded_model_udf qui charge un modèle MLflow entraîné sur des données de risque de prêt. Les colonnes de données utilisées pour faire la prédiction sont passées comme argument à la UDF. La table loan_risk_predictions calcule les prédictions pour chaque ligne de loan_risk_input_data.
%pip install mlflow==2.20.2
from pyspark import pipelines as dp
import mlflow
from pyspark.sql.functions import struct
run_id = "mlflow_run_id"
model_name = "the_model_name_in_run"
model_uri = f"runs:/{run_id}/{model_name}"
loaded_model_udf = mlflow.pyfunc.spark_udf(spark, model_uri=model_uri)
categoricals = ["term", "home_ownership", "purpose",
"addr_state","verification_status","application_type"]
numerics = ["loan_amnt", "emp_length", "annual_inc", "dti", "delinq_2yrs",
"revol_util", "total_acc", "credit_length_in_years"]
features = categoricals + numerics
@dp.materialized_view(
comment="GBT ML predictions of loan risk",
table_properties={
"quality": "gold"
}
)
def loan_risk_predictions():
return spark.read.table("loan_risk_input_data")
.withColumn('predictions', loaded_model_udf(struct(features)))
Conserver les suppressions ou mises à jour manuelles
Les pipelines vous permettent de supprimer ou de mettre à jour manuellement des enregistrements d'une table et d'effectuer une opération de refresh pour recalculer les tables en aval.
Par default, les pipelines recalculent les résultats des tables en fonction des données d'entrée chaque fois qu'elles sont mises à jour. Vous devez donc vous assurer que l'enregistrement supprimé n'est pas rechargé à partir des données sources. Le fait de définir la propriété de table pipelines.reset.allowed sur false empêche les refresh d'une table, mais n'empêche pas les écritures incrémentielles dans les tables ou l'afflux de nouvelles données dans la table.
Le diagramme suivant illustre un exemple utilisant deux tables de streaming :
raw_user_tableingère des données utilisateur brutes à partir d’une source.bmi_tablecalcule les scores BMI de manière incrémentielle en utilisant le poids et la taille deraw_user_table.
Vous souhaitez supprimer ou mettre à jour manuellement des enregistrements d'utilisateur à partir du raw_user_table et recalculer le bmi_table.

Le code suivant montre comment définir la propriété de table pipelines.reset.allowed sur false afin de désactiver le refresh complet pour raw_user_table pour que les modifications prévues soient conservées au fil du temps, mais que les tables en aval soient recalculées lors de l'exécution d'une mise à jour de pipeline :
CREATE OR REFRESH STREAMING TABLE raw_user_table
TBLPROPERTIES(pipelines.reset.allowed = false)
AS SELECT * FROM STREAM read_files("/databricks-datasets/iot-stream/data-user", format => "csv");
CREATE OR REFRESH STREAMING TABLE bmi_table
AS SELECT userid, (weight/2.2) / pow(height*0.0254,2) AS bmi FROM STREAM(raw_user_table);