Aller au contenu principal

Comparer Spark Connect à Spark Classic

Spark Connect est un protocole basé sur gRPC au sein d’Apache Spark qui spécifie comment une application cliente peut communiquer avec un serveur Spark distant. Il permet l’exécution à distance de charges de travail Spark à l’aide de l’API DataFrame.

Spark Connect est utilisé dans les cas suivants :

  • Notebooks Scala avec Databricks Runtime version 13,3 et versions ultérieures, sur un compute standard
  • Notebooks Python avec Databricks Runtime version 14.3 et supérieur, sur un compute standard
  • Compute serverless
  • Databricks Connect
  • LakeFlow Pipelines avec des versions d'environnement configurées

Bien que Spark Connect et Spark Classic utilisent tous deux l'exécution paresseuse pour les Transformations, il existe d'importantes différences à connaître pour éviter les comportements inattendus et les problèmes de performances lors de la migration du code existant de Spark Classic vers Spark Connect ou lors de l'écriture de code qui doit fonctionner avec les deux.

Lazy vs eager

La principale différence entre Spark Connect et Spark Classic est que Spark Connect reporte l’analyse et la résolution de noms au moment de l’exécution, comme résumé dans le tableau suivant.

Aspect

Spark Classic

Spark Connect

Exécution de la query

Lazy

Lazy

Analyse de schéma

Anticipé

Lazy

Accès au schéma

Local

Déclenche l'RPC et met en cache le schéma lors du premier accès

Vues temporaires

Plan intégré

Recherche par nom

Sérialisation d'UDF

À la création

Lors de l'exécution

Aspect

Spark Classic

Spark Connect

Exécution de la query

Lazy

Lazy

Analyse de schéma

Anticipé

Lazy

Accès au schéma

Local

Déclenche l'RPC et met en cache le schéma lors du premier accès

Vues temporaires

Plan intégré

Recherche par nom

Sérialisation d'UDF

À la création

Lors de l'exécution

Exécution de la query

Spark Classic et Spark Connect suivent tous deux le même modèle d'exécution paresseuse pour l'exécution des requêtes.

Dans Spark Classic, les transformations de DataFrame (telles que filter et limit) sont paresseuses. Cela signifie qu'elles ne sont pas exécutées immédiatement, mais sont encodées dans un plan logique. Le calcul réel n'est déclenché que par une action (telle que show(), collect()).

Spark Connect suit un modèle d'évaluation différée similaire. Les Transformations sont élaborées côté client et envoyées sous forme de plans non résolus au serveur. Le serveur effectue ensuite l'analyse et l'exécution nécessaires lorsqu'une action est appelée.

Aspect

Spark Classic

Spark Connect

Transformations : df.filter(...), df.select(...), df.limit(...)

Exécution paresseuse

Exécution paresseuse

query SQL : spark.sql("select …")

Exécution paresseuse

Exécution paresseuse

Actions : df.collect(), df.show()

Exécution immédiate

Exécution immédiate

Commandes SQL : spark.sql("insert …"), spark.sql("create …")

Exécution immédiate

Exécution immédiate

Aspect

Spark Classic

Spark Connect

Transformations : df.filter(...), df.select(...), df.limit(...)

Exécution paresseuse

Exécution paresseuse

query SQL : spark.sql("select …")

Exécution paresseuse

Exécution paresseuse

Actions : df.collect(), df.show()

Exécution immédiate

Exécution immédiate

Commandes SQL : spark.sql("insert …"), spark.sql("create …")

Exécution immédiate

Exécution immédiate

Analyse de schéma

Spark Classic effectue l'analyse de manière proactive pendant la construction du plan logique. Cette phase d'analyse convertit le plan non résolu en un plan logique entièrement résolu et vérifie que l'opération peut être exécutée par Spark. L'un des principaux avantages de réaliser ce travail de manière proactive est que les utilisateurs reçoivent un retour immédiat lorsqu'une erreur est commise. Par exemple, l'exécution de spark.sql("select 1 as a, 2 as b").filter("c > 1") générera une erreur prématurément, indiquant que la colonne c est introuvable.

Spark Connect diffère de Classic car le client construit des plans non résolus pendant la transformation et diffère leur analyse. Toute opération qui nécessite un plan résolu — comme l'accès à un schéma, l'explication du plan, la persistance d'un DataFrame ou l'exécution d'une action — amène le client à envoyer les plans non résolus au serveur via RPC. Le serveur effectue ensuite une analyse complète pour obtenir son plan logique résolu et réaliser l'opération. Par exemple, spark.sql("select 1 as a, 2 as b").filter("c > 1") ne générera pas d'erreur car le plan non résolu est uniquement côté client, mais sur df.columns ou df.show(), une erreur sera générée car le plan non résolu est envoyé au serveur pour analyse.

Contrairement à l'exécution de query, Spark Classic et Spark Connect diffèrent quant au moment où l'analyse du schéma a lieu.

Aspect

Spark Classic

Spark Connect

Transformations : df.filter(...), df.select(...), df.limit(...)

Anticipé

Lazy

Accès au Schéma : df.columns, df.schema, df.isStreaming

Anticipé

Anticipé

Trigger une requête RPC d'analyse, contrairement à Spark Classic.

Actions : df.collect(), df.show()

Anticipé

Anticipé

État de session dépendant des DataFrames : UDFs, vues temporaires, configs

Anticipé

Lazy

Évalué pendant l'exécution du plan du DataFrame

État de session dépendant des vues temporaires : Fonctions définies par l'utilisateur (UDF), autres vues temporaires, configurations

Anticipé

Anticipé

L'analyse est Trigger hâtivement lors de la création de la vue temporaire.

Aspect

Spark Classic

Spark Connect

Transformations : df.filter(...), df.select(...), df.limit(...)

Anticipé

Lazy

Accès au Schéma : df.columns, df.schema, df.isStreaming

Anticipé

Anticipé

Trigger une requête RPC d'analyse, contrairement à Spark Classic.

Actions : df.collect(), df.show()

Anticipé

Anticipé

État de session dépendant des DataFrames : UDFs, vues temporaires, configs

Anticipé

Lazy

Évalué pendant l'exécution du plan du DataFrame

État de session dépendant des vues temporaires : Fonctions définies par l'utilisateur (UDF), autres vues temporaires, configurations

Anticipé

Anticipé

L'analyse est Trigger hâtivement lors de la création de la vue temporaire.

Bonnes pratiques

La différence entre l'analyse paresseuse et l'analyse diligente signifie qu'il existe des bonnes pratiques à suivre pour éviter les comportements inattendus et les problèmes de performance, spécifiquement ceux causés par l'écrasement des noms de vues temporaires, la capture de variables externes dans les UDF, la détection retardée des erreurs et l'accès excessif au schéma sur les nouveaux DataFrames.

Créer des noms de vues temporaires uniques

Dans Spark Connect, le DataFrame ne stocke qu'une référence à la vue temporaire par son nom. En conséquence, si la vue temporaire est remplacée ultérieurement, les données du DataFrame changeront également car il recherche la vue par nom au moment de l'exécution.

Ce comportement diffère de Spark Classic, où le plan logique de la vue temporaire est intégré au plan du DataFrame au moment de sa création. Tout remplacement ultérieur de la vue temporaire n’affecte pas le cadre de données original.

Pour atténuer la différence, créez toujours des noms de vues temporaires uniques. Par exemple, incluez un UUID dans le nom de la vue. Ceci évite d'affecter les DataFrames existants qui font référence à une vue temporaire précédemment enregistrée.

Python
import uuid
def create_temp_view_and_create_dataframe(x):
temp_view_name = f"`temp_view_{uuid.uuid4()}`" # Use a random name to avoid conflicts.
spark.range(x).createOrReplaceTempView(temp_view_name)
return spark.table(temp_view_name)

df10 = create_temp_view_and_create_dataframe(10)
assert len(df10.collect()) == 10

df100 = create_temp_view_and_create_dataframe(100)
assert len(df10.collect()) == 10 # It works as expected now.
assert len(df100.collect()) == 100

Encapsuler les définitions d'UDF

Il est généralement considéré comme une mauvaise pratique pour les UDF de dépendre de variables externes mutables, car cela introduit des dépendances implicites, peut entraîner un comportement non déterministe et réduit la composabilité. Toutefois, si vous avez un tel modèle, soyez conscient de l'écueil suivant :

Dans Spark Connect, les UDF Python sont paresseuses. Leur sérialisation et leur enregistrement sont différés jusqu'à l'heure d'exécution. Dans l’exemple suivant, l’UDF est seulement sérialisée et upload au cluster Spark pour exécution lorsque show() est appelée.

Python
from pyspark.sql.functions import udf

x = 123

@udf("INT")
def foo():
return x


df = spark.range(1).select(foo())
x = 456
df.show() # Prints 456

Ce comportement diffère de Spark Classic, où les UDF sont créées de manière anticipée. Dans Spark Classic, la valeur de x est capturée au moment de la création de l'UDF, de sorte que les modifications ultérieures apportées à x n'affectent pas l'UDF déjà créée.

Si vous devez modifier la valeur des variables externes dont dépend une UDF, utilisez une fabrique de fonctions (fermeture avec liaison précoce) pour capturer correctement les valeurs des variables. Plus précisément, encapsulez la création de l'UDF dans une fonction d'aide pour capturer la valeur d'une variable dépendante.

Python
from pyspark.sql.functions import udf

def make_udf(value):
def foo():
return value
return udf(foo)


x = 123
foo_udf = make_udf(x)
x = 456
df = spark.range(1).select(foo_udf())
df.show() # Prints 123 as expected

En encapsulant la définition d'UDF à l'intérieur d'une autre fonction (make_udf), nous créons une nouvelle portée où la valeur actuelle de x est passée en tant qu'argument. Cela garantit que chaque UDF générée dispose de sa propre copie du champ, liée au moment où l'UDF est créée.

Trigger l'analyse anticipée pour la détection d'erreurs

La gestion des erreurs suivante est utile dans Spark Classic car elle effectue une analyse immédiate, ce qui permet de lever les exceptions rapidement. Cependant, dans Spark Connect, ce code ne pose aucun problème, car il ne construit qu'un plan local non résolu sans Trigger d'analyse.

Python
df = spark.createDataFrame([("Alice", 25), ("Bob", 30)], ["name", "age"])

try:
df = df.select("name", "age")
df = df.withColumn(
"age_group",
when(col("age") < 18, "minor").otherwise("adult"))
df = df.filter(col("age_with_typo") > 6) # The use of non-existing column name will not throw analysis exception in Spark Connect
except Exception as e:
print(f"Error: {repr(e)}")

Si votre code dépend de l'exception d'analyse et que vous voulez la rattraper, vous pouvez trigger une analyse anticipée, par exemple avec df.columns, df.schema ou df.collect().

Python
try:
df = ...
df.columns # This will trigger eager analysis
except Exception as e:
print(f"Error: {repr(e)}")

Évitez trop de requêtes d'analyse anticipées

Les performances peuvent être améliorées si vous évitez les requêtes d'analyse sur un grand nombre de DataFrames.

Création de nouveaux DataFrames étape par étape et accès à leur schéma à chaque itération

Lorsque vous créez un grand nombre de nouveaux DataFrames, évitez l'utilisation excessive d'appels déclenchant une analyse immédiate sur ceux-ci (tels que df.columns, df.schema). Vous pouvez accéder au schéma du même DataFrame plusieurs fois, mais déclencher une analyse sur de nombreux DataFrames nouvellement créés aura un impact sur les performances.

Par exemple, lorsque vous ajoutez itérativement des colonnes à un DataFrame dans une boucle et vérifiez si chaque colonne existe déjà avant de l’ajouter, appeler df.columns sur chaque nouveau DataFrame Trigger une requête d’analyse à chaque itération. Pour éviter cela, conservez un ensemble pour suivre les noms de colonnes au lieu d’accéder à plusieurs reprises au schéma du DataFrame.

Python
df = spark.range(10)
columns = set(df.columns) # Maintain the set of column names
for i in range(200):
new_column_name = str(i)
# if new_column_name not in df.columns: # Bad practice. The `df.columns` call causes an analysis request on the newly created DataFrame in every iteration.
if new_column_name not in columns: # Check the set without triggering analysis
df = df.withColumn(new_column_name, F.col("id") + i)
columns.add(new_column_name)
df.show()

Évitez d'accéder aux schémas pour un grand nombre de DataFrames intermédiaires

Un autre cas similaire est la création d'un grand nombre de DataFrames intermédiaires inutiles et leur analyse. Dans le cas suivant, pour extraire les noms de champ de chaque colonne d'un type structuré, obtenez directement les informations de champ StructType du schéma du DataFrame au lieu de créer des DataFrames intermédiaires.

Python
from pyspark.sql.types import StructType

df = ...
struct_column_fields = {
# column_schema.name: df.select(column_schema.name + ".*").columns # Bad practice. This creates an intermediate DataFrame and triggers an analysis request for each StructType column.
column_schema.name: [f.name for f in column_schema.dataType.fields] # Access StructType fields directly from the schema, avoiding analysis on intermediate DataFrames.
for column_schema in df.schema
if isinstance(column_schema.dataType, StructType)
}
print(struct_column_fields)