Compatibilité des versions d'environnement
Bêta
Les versions de l’environnement pour Lakeflow Pipelines sont en Bêta.
Les pipelines avec une version d'environnement définie exécutent du code Python via Spark Connect. Cette page couvre ce qui est incompatible, ce qui se comporte différemment, comment analyser un pipeline pour détecter les modèles affectés, et comment migrer un pipeline existant.
Limitations
Les versions d'environnement ne sont pas encore compatibles avec toutes les fonctionnalités du pipeline. Une exécution de pipeline avec un ensemble de versions d'environnement échoue si le code Python du pipeline effectue l'une des opérations suivantes :
- Modifie l'état de la session Spark au sein d'une fonction décorée par un décorateur de pipelines. Exemples :
spark.conf.set(...),spark.sql("USE CATALOG ...")etcreateOrReplaceTempView. - Utilise les APIs PySpark qui ne sont pas disponibles dans Spark Connect, notamment
SparkContext,RDD,SQLContextet toutes les APIs Py4J. Consultez ce qui est pris en charge dans Spark Connect.
Si l'activation d'une version d'environnement sur un pipeline entraîne son échec, la désactivation de la version d'environnement ramène le pipeline à son état précédent.
Changements de comportement
Spark Connect présente un petit nombre de différences de comportement par rapport au runtime PySpark classique. Voir Spark Connect vs. Spark classique pour la référence complète. La recherche de compatibilité détecte ces modèles à l'avance et bloque l'activation jusqu'à ce qu'ils soient résolus, afin que vous puissiez les trouver et les corriger avant qu'ils n'affectent les données de production.
Dans un pipeline, les situations les plus courantes où le comportement peut différer sont :
- Entrelacement de la construction de DataFrame et de la mutation de session
- UDF qui référencent un état Python mutable
Construction de DataFrame entrelacées et mutation de session
Lorsqu'un pipeline construit un DataFrame, puis modifie l'état de la session Spark (par exemple, modifie le catalogue ou le schéma default, définit une configuration, remplace une vue temporaire ou ré-enregistre une UDF), puis utilise le DataFrame :
- Sans version d'environnement, le DataFrame utilise l'état de session de pré-mutation .
- Avec une version d'environnement, le DataFrame utilise l'état de session post-mutation .
Par exemple :
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Sans version d'environnement, mytable contient [(1, "Original Row")]. Avec une version d'environnement, mytable contient [(2, "Replaced Row")].
UDFs qui référencent un état Python mutable
Lorsqu'une UDF fait référence à une variable globale Python dont la valeur change après la définition de l'UDF :
- Sans version d’environnement, l’UDF utilise la valeur **la plus récente** de la variable.
- Avec une version d'environnement, l'UDF utilise la valeur au moment où l'UDF a été définie .
Par exemple :
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Sans version d'environnement, my_mv contient [("alex_b",)]. Avec une version d'environnement, my_mv contient [("alex_a",)].
Si un pipeline repose sur l'un ou l'autre de ces modèles, auditez-le avant d'activer une version d'environnement.
Analyse de compatibilité
Le scan de compatibilité vous aide à trouver des modèles de code dans votre pipeline qui produiraient des résultats différents sous une version d'environnement, avant que vous n'en activiez une. Le scan est sur option. Lorsque le scan est activé sur un pipeline :
- Chaque exécution de pipeline émet un événement
BehaviorChangeInSparkConnectWARNdans le log des événements du pipeline par modèle détecté. - Vous ne pouvez pas activer une version d'environnement sur le pipeline tant que vous n'avez pas résolu tous les avertissements de compatibilité de la mise à jour réussie précédente.
Si le scan n'est pas activé, aucun événement n'est émis et l'activation de environment_version n'est pas bloquée. Databricks recommande d'activer l'analyse et de résoudre les modèles détectés avant d'activer une version d'environnement sur le pipeline.
Activer l'analyse sur un pipeline
Vous pouvez activer l'analyse de compatibilité en ajoutant la configuration du pipeline pipelines.environmentVersion.enableCompatibilityScan. Vous pouvez ajouter la configuration via l’interface utilisateur de l’éditeur de pipeline ou en ajoutant une entrée au JSON de configuration du pipeline.
Via l'interface utilisateur :
- Dans l'éditeur de pipeline, cliquez sur Paramètres .
- Trouvez la section Configuration dans les paramètres du pipeline.
- Cliquez sur
Ajouter une configuration .
- Saisissez
pipelines.environmentVersion.enableCompatibilityScancomme clé ettruecomme valeur. - Enregistrez les paramètres du pipeline.
Dans le JSON du pipeline :
Ajoutez l'entrée suivante au bloc configuration :
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Workflow recommandé
- Activez l'analyse sur le pipeline.
- Trigger une exécution de pipeline.
- query the Logs des événements du pipeline pour les
BehaviorChangeInSparkConnectévénementsWARN. Consultez la référence des événements de compatibilité pour la liste complète des codes de problème, des modèles d'exemple et des correctifs suggérés. - Mettez à jour le code du pipeline pour supprimer les modèles détectés et exécutez le pipeline à nouveau jusqu'à ce qu'aucun autre événement ne soit émis.
- Ajoutez
environment_versionau pipeline à l'aide de l'une des méthodes décrites dans Activer une version d'environnement sur un pipeline.
Si vous estimez qu'un avertissement de compatibilité est un faux positif et que vous souhaitez activer environment_version quand même, supprimez l'entrée pipelines.environmentVersion.enableCompatibilityScan de la configuration du pipeline pour contourner la vérification. (Il n'est pas permis de définir la valeur sur false — vous devez supprimer l'entrée entièrement.)
La vérification préliminaire ne s'exécute pas sur les pipelines qui n'ont pas de mise à jour précédente, ou sur les pipelines qui ont déjà une version d'environnement définie.
Migrer un pipeline existant vers des versions d'environnement
Pour migrer un pipeline existant qui n'utilise pas encore de version d'environnement, suivez ce workflow de bout en bout. Il vous guide à travers la recherche de modèles de code susceptibles de se comporter différemment avec Spark Connect, leur correction et le déploiement sécurisé de la version de l'environnement.
-
**Activez l'analyse de compatibilité sur le pipeline.** Activez l'analyse sur le pipeline comme décrit dans Analyse de compatibilité. C'est ce qui fait apparaître les modèles détectés dans le journal d'événements et ce qui permet la vérification préalable qui protège votre tentative d'activation.
-
Trigger a pipeline run and review compatibility events. Trigger une mise à jour normale du pipeline. Une fois l'opération terminée, interrogez le Logs d'événements du pipeline pour
BehaviorChangeInSparkConnectWARNévénements. Chaque événement rapporte un modèle détecté. Consultez la documentation de référence sur les événements de compatibilité pour obtenir la liste complète des codes de problème, des modèles d'exemple et des correctifs suggérés. -
Mettez à jour votre code de pipeline pour résoudre les modèles détectés. Pour chaque modèle détecté, mettez à jour votre code de pipeline en suivant la correction suggérée. Après chaque modification, Trigger une autre mise à jour du pipeline et vérifiez que les événements correspondants n'apparaissent plus. Répétez l'opération jusqu'à ce que le Logs ne signale plus aucun événement de compatibilité pour une mise à jour réussie.
-
**Activez la version d'environnement sur le pipeline.** Une fois que la mise à jour réussie la plus récente n’a pas d’événements de compatibilité, ajoutez
environment_versionau pipeline à l'aide de l'interface utilisateur, de l'API ou du bundle, comme décrit dans Activer une version d'environnement sur un pipeline. La prochaine mise à jour s'exécute avec Spark Connect et la version de langage Python épinglée et les bibliothèques préinstallées.Si la mise à jour échoue parce que des avertissements de compatibilité existent toujours, supprimez le
environment_version, revenez à l'étape 2 et résolvez les avertissements restants avant de réessayer. -
Vérifiez la migration. Une fois la première mise à jour avec la version de l'environnement terminée, vérifier :
- L'événement
create_updatedans l'event Logs montreenvironment_versiondéfini sur la valeur attendue. - Le pipeline produit les données attendues et aucun nouvel événement d'erreur n'apparaît.
- Vérifiez les tables en aval pour détecter toute différence de comportement subtile décrite dans Modifications de comportement.
- L'événement
Retour
Si le pipeline se comporte mal après la migration, supprimez le environment_version des paramètres du pipeline. La prochaine mise à jour s'exécute avec la configuration d'exécution Python précédente. Utilisez l'exécution restaurée pour déboguer, puis répétez la migration à partir de l'étape 2 après avoir identifié et résolu le problème.
Référence des événements de compatibilité
Lorsqu’une analyse de compatibilité est activée sur un pipeline, elle émet un événement BehaviorChangeInSparkConnect WARN dans le log des événements du pipeline par modèle détecté. Lorsque l'analyse est activée et que la mise à jour précédente réussie a détecté des modèles, le pipeline bloque également l'activation de environment_version jusqu'à ce que les modèles soient traités.
Chaque événement signale un code de problème unique qui identifie ce qui a été détecté. Pour rechercher un code, trouvez-le dans la table Codes de problème — chaque ligne renvoie à la section de catégorie qui contient un modèle d'exemple et la correction suggérée.
Forme d'événement
BehaviorChangeInSparkConnect les événements suivent le schéma standard du pipeline event log schema:
event_typeestbehavior_change_in_spark_connect.levelestWARN.detailscontient l'objetbehavior_change_in_spark_connect, qui a un seul champissue. La valeur du problème est l'un des codes listés ci-dessous.messageest une description lisible par l'homme du modèle détecté.
Codes de problème
Catégorie | Code du problème | Description |
|---|---|---|
| Le catalogue default a été modifié après la création d'un DataFrame. Le DataFrame existant peut résoudre des tables en utilisant le nouveau catalogue default. | |
|
| |
| La base de données default a été modifiée après la création d'un DataFrame. Le DataFrame existant peut résoudre des tables en utilisant la nouvelle base de données default. | |
|
| |
| La fonction de flux appelle une commande de point de contrôle. | |
| La fonction de flux crée de manière anticipée une vue DataFrame ( | |
| La fonction de flux crée un profil de ressource. | |
| La fonction de flux appelle | |
| La fonction de flux effectue une | |
| La fonction de flux effectue une opération Spark ML impatiente. | |
| La fonction de flux enregistre une source de données Python. | |
| La fonction de flux opère sur un handle de query de streaming actif. | |
| La fonction de flux enregistre ou supprime un écouteur de query en streaming. | |
| La fonction de flux appelle | |
| La fonction de flux effectue une opération | |
| La fonction de flux effectue une opération | |
| La fonction de flux démarre une requête de streaming ( | |
|
| |
|
| |
|
| |
| Une vue temporaire globale a été remplacée après la création d'un DataFrame la référençant. Le remplacement peut être reflété dans le DataFrame existant. | |
| Une vue temporaire a été remplacée après la création d'un DataFrame la référençant. Le remplacement peut être reflété dans le DataFrame existant. | |
| Un UDF a été réenregistré sous le même nom après la création d'un DataFrame le référençant. Le DataFrame existant peut utiliser la nouvelle définition d'UDF. | |
| Une UDTF a été réenregistrée avec le même nom après la création d’un DataFrame y faisant référence. Le DataFrame existant peut utiliser la nouvelle définition d'UDTF. | |
| Une UDF fait référence à une variable Python globale modifiable. Avec une version d'environnement, l'UDF utilise la valeur de la variable au moment où l'UDF a été définie, et non au moment de l'invocation. | |
| Une UDTF référence une variable Python globale et mutable. Avec une version d'environnement, l'UDTF utilise la valeur de la variable au moment où l'UDTF a été définie, et non au moment de l'invocation. |
Mutations de bases de données et de catalogues
Ces problèmes se produisent lorsque le code de pipeline modifie la base de données ou le catalogue default. Avec une version d'environnement, les DataFrames construits avant la mutation peuvent résoudre des tables à l'aide de la nouvelle base de données ou du nouveau catalogue.
Exemple de modèle qui déclenche un événement :
from pyspark import pipelines as dp
spark.sql("USE CATALOG marketing")
df = spark.read.table("events")
spark.sql("USE CATALOG sales") # changes the default catalog after df was created
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Sans version d'environnement, df résout events à partir du catalogue marketing. Avec une version d'environnement, df résout events à partir du catalogue sales.
Correction suggérée : qualifiez entièrement les noms de table afin que la résolution ne dépende pas du catalogue ou de la base de données par default, et évitez de modifier le catalogue ou la base de données par default entre la création et l'utilisation du DataFrame.
from pyspark import pipelines as dp
df = spark.read.table("marketing.default.events")
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Mutations de la configuration Spark
Ces problèmes sont émis lorsque le code de pipeline modifie la configuration Spark de manières qui peuvent changer le comportement du DataFrame sous une version d'environnement.
Exemple de modèle qui déclenche un événement :
from pyspark import pipelines as dp
df = spark.read.table("events")
spark.conf.set("spark.sql.ansi.enabled", "true") # changes session conf after df was created
@dp.materialized_view
def events_strict():
return df.selectExpr("CAST(price AS INT) AS price")
Sans version d’environnement, la conversion utilise la valeur de conf au moment de la création du DataFrame. Avec une version d’environnement, la conversion utilise spark.sql.ansi.enabled=true et peut échouer en cas d’entrée non valide.
**Correction suggérée :** Configurez toutes les configurations Spark requises en haut du fichier de pipeline, avant la création de tout DataFrame. Pour la configuration par query, utilisez le paramètre configuration du pipeline dans la spécification du pipeline.
Remplacements de vue temporaire
Ces problèmes sont émis lorsque le code de pipeline remplace une vue temporaire après la création d'un DataFrame la référençant. Avec une version d'environnement, le DataFrame existant peut refléter le nouveau contenu de la vue.
Exemple de modèle qui déclenche un événement :
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Sans version d'environnement, mytable contient [(1, "Original Row")]. Avec une version d'environnement, mytable contient [(2, "Replaced Row")].
Correction suggérée : Créez chaque vue temporaire une seule fois et ne la remplacez pas. Si vous avez besoin de plusieurs vues avec des données associées, donnez à chacune un nom distinct.
Mutations UDF et UDTF
Ces problèmes surviennent lorsque le code du pipeline modifie une UDF ou une UDTF de manière à altérer le comportement sous une version d'environnement.
Exemple de modèle qui déclenche un événement :
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Sans version d'environnement, my_mv contient [("alex_b",)]. Avec une version d'environnement, my_mv contient [("alex_a",)].
Correction suggérée : Transmettez les valeurs à l'UDF en tant qu'arguments au lieu de les capturer à partir de variables globales Python, ou définissez la variable globale avant de définir l'UDF et ne la modifiez pas par la suite.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, lit, udf
@udf
def append_suffix(s, suffix):
return s + suffix
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(append_suffix(col("name"), lit("b")))
Exécution immédiate au sein des fonctions de flux
Ces problèmes sont émis lorsque le code du pipeline exécute une commande Spark anticipée à l'intérieur d'une fonction décorée par un décorateur de pipelines (@table, @materialized_view, etc.). Les fonctions de flux sont censées définir et retourner un DataFrame ; les commandes anticipées qui écrivent des données, gèrent des queries de streaming, enregistrent des ressources ou exécutent des Opérations ML ne sont pas autorisées à l'intérieur d'une fonction de flux avec une version d'environnement définie.
Correctif suggéré : Déplacez l'opération hâtive en dehors de la fonction de flux et renvoyez un DataFrame à partir de la fonction de flux. Les effets secondaires, tels que l'écriture dans une table ou le démarrage d'une query de streaming, doivent se trouver en dehors de la définition du pipeline ; le moteur de pipeline gère la matérialisation du DataFrame renvoyé par la fonction de flux.
Rechercher les événements de compatibilité dans le journal des Logs
La query suivante renvoie tous les événements de compatibilité pour un pipeline, classés du plus récent au plus ancien :
SELECT
timestamp,
message,
details:behavior_change_in_spark_connect:issue AS issue
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
ORDER BY timestamp DESC;
Pour compter les événements par code de problème parmi les mises à jour récentes :
SELECT
details:behavior_change_in_spark_connect:issue AS issue,
COUNT(*) AS occurrences
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
GROUP BY 1
ORDER BY occurrences DESC;
Pour savoir comment interroger le event log, consultez Interroger le event log.
Ressources supplémentaires
- Configurer les versions d'environnement pour les pipelines — aperçu des fonctionnalités, comment activer une version d'environnement.
- Schéma du log d'événements du pipeline — schéma complet du log d'événements du pipeline.
- Journal des événements du pipeline — comment query le journal des événements du pipeline.