Aller au contenu principal

Compatibilité des versions d'environnement

info

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 ...") et createOrReplaceTempView.
  • Utilise les APIs PySpark qui ne sont pas disponibles dans Spark Connect, notamment SparkContext, RDD, SQLContext et 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 :

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 :

Python
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 :

Python
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 BehaviorChangeInSparkConnect WARN dans 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 :

  1. Dans l'éditeur de pipeline, cliquez sur Paramètres .
  2. Trouvez la section Configuration dans les paramètres du pipeline.
  3. Cliquez sur Icône Plus. Ajouter une configuration .
  4. Saisissez pipelines.environmentVersion.enableCompatibilityScan comme clé et true comme valeur.
  5. Enregistrez les paramètres du pipeline.

Dans le JSON du pipeline :

Ajoutez l'entrée suivante au bloc configuration :

JSON
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}

Workflow recommandé

  1. Activez l'analyse sur le pipeline.
  2. Trigger une exécution de pipeline.
  3. query the Logs des événements du pipeline pour les BehaviorChangeInSparkConnect événements WARN. 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.
  4. 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.
  5. Ajoutez environment_version au 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.

  1. **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.

  2. 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 BehaviorChangeInSparkConnect WARN é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.

  3. 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.

  4. **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_version au 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.

  5. 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_update dans l'event Logs montre environment_version dé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.

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_type est behavior_change_in_spark_connect.
  • level est WARN.
  • details contient l'objet behavior_change_in_spark_connect, qui a un seul champ issue. La valeur du problème est l'un des codes listés ci-dessous.
  • message est une description lisible par l'homme du modèle détecté.

Codes de problème

Catégorie

Code du problème

Description

Mutations de bases de données et de catalogues

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations de bases de données et de catalogues

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE CATALOG a été appelée en dehors d'une fonction décorée par un décorateur de pipelines. Le catalogue default peut changer de manière inattendue pour les Opérations suivantes.

Mutations de bases de données et de catalogues

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations de bases de données et de catalogues

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE DATABASE a été appelée en dehors d'une fonction décorée par un décorateur de pipelines. La base de données default peut changer de manière inattendue pour les opérations ultérieures.

Exécution immédiate au sein des fonctions de flux

CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle une commande de point de contrôle.

Exécution immédiate au sein des fonctions de flux

CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux crée de manière anticipée une vue DataFrame (createOrReplaceTempView ou similaire).

Exécution immédiate au sein des fonctions de flux

CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux crée un profil de ressource.

Exécution immédiate au sein des fonctions de flux

GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle spark.resources ou une API de Ressource associée.

Exécution immédiate au sein des fonctions de flux

MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une MERGE INTO anticipée sur une table cible.

Exécution immédiate au sein des fonctions de flux

ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération Spark ML impatiente.

Exécution immédiate au sein des fonctions de flux

REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux enregistre une source de données Python.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux opère sur un handle de query de streaming actif.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux enregistre ou supprime un écouteur de query en streaming.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle spark.streams pour gérer les requêtes de streaming.

Exécution immédiate au sein des fonctions de flux

WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération DataFrameWriterV2 anticipée.

Exécution immédiate au sein des fonctions de flux

WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération DataFrame.write anticipée.

Exécution immédiate au sein des fonctions de flux

WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux démarre une requête de streaming (writeStream.start()).

Modifications de la configuration Spark

CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED

spark.conf.set() ou spark.conf.unset() a été appelé à l’intérieur d’une fonction décorée par un décorateur de pipelines. Ceci n'est pas pris en charge avec une version d'environnement.

Modifications de la configuration Spark

SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.set() a été appelée en dehors d'une fonction décorée par un décorateur de pipelines après la création d'un DataFrame. La modification de configuration peut affecter le DataFrame existant au moment de l'exécution.

Modifications de la configuration Spark

UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.unset() a été appelée en dehors d'une fonction décorée par un décorateur de pipelines après la création d'un DataFrame. La modification de configuration peut affecter le DataFrame existant au moment de l'exécution.

Remplacements de vues temporaires

REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Remplacements de vues temporaires

REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Catégorie

Code du problème

Description

Mutations de bases de données et de catalogues

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations de bases de données et de catalogues

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE CATALOG a été appelée en dehors d'une fonction décorée par un décorateur de pipelines. Le catalogue default peut changer de manière inattendue pour les Opérations suivantes.

Mutations de bases de données et de catalogues

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations de bases de données et de catalogues

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE DATABASE a été appelée en dehors d'une fonction décorée par un décorateur de pipelines. La base de données default peut changer de manière inattendue pour les opérations ultérieures.

Exécution immédiate au sein des fonctions de flux

CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle une commande de point de contrôle.

Exécution immédiate au sein des fonctions de flux

CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux crée de manière anticipée une vue DataFrame (createOrReplaceTempView ou similaire).

Exécution immédiate au sein des fonctions de flux

CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux crée un profil de ressource.

Exécution immédiate au sein des fonctions de flux

GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle spark.resources ou une API de Ressource associée.

Exécution immédiate au sein des fonctions de flux

MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une MERGE INTO anticipée sur une table cible.

Exécution immédiate au sein des fonctions de flux

ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération Spark ML impatiente.

Exécution immédiate au sein des fonctions de flux

REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux enregistre une source de données Python.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux opère sur un handle de query de streaming actif.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux enregistre ou supprime un écouteur de query en streaming.

Exécution immédiate au sein des fonctions de flux

STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux appelle spark.streams pour gérer les requêtes de streaming.

Exécution immédiate au sein des fonctions de flux

WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération DataFrameWriterV2 anticipée.

Exécution immédiate au sein des fonctions de flux

WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux effectue une opération DataFrame.write anticipée.

Exécution immédiate au sein des fonctions de flux

WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

La fonction de flux démarre une requête de streaming (writeStream.start()).

Modifications de la configuration Spark

CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED

spark.conf.set() ou spark.conf.unset() a été appelé à l’intérieur d’une fonction décorée par un décorateur de pipelines. Ceci n'est pas pris en charge avec une version d'environnement.

Modifications de la configuration Spark

SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.set() a été appelée en dehors d'une fonction décorée par un décorateur de pipelines après la création d'un DataFrame. La modification de configuration peut affecter le DataFrame existant au moment de l'exécution.

Modifications de la configuration Spark

UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.unset() a été appelée en dehors d'une fonction décorée par un décorateur de pipelines après la création d'un DataFrame. La modification de configuration peut affecter le DataFrame existant au moment de l'exécution.

Remplacements de vues temporaires

REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Remplacements de vues temporaires

REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Mutations UDF et UDTF

UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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 :

Python
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.

Python
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 :

Python
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 :

Python
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 :

Python
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.

Python
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 :

SQL
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 :

SQL
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