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 |
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 : | Exécution paresseuse | Exécution paresseuse |
query SQL : | Exécution paresseuse | Exécution paresseuse |
Actions : | Exécution immédiate | Exécution immédiate |
Commandes SQL : | 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 : | Anticipé | Lazy |
Accès au Schéma : | Anticipé | Anticipé Trigger une requête RPC d'analyse, contrairement à Spark Classic. |
Actions : | 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
- Scala
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
import java.util.UUID
def createTempViewAndDataFrame(x: Int) = {
val tempViewName = s"`temp_view_${UUID.randomUUID()}`"
spark.range(x).createOrReplaceTempView(tempViewName)
spark.table(tempViewName)
}
val df10 = createTempViewAndDataFrame(10)
assert(df10.collect().length == 10)
val df100 = createTempViewAndDataFrame(100)
assert(df10.collect().length == 10) // Works as expected
assert(df100.collect().length == 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.
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
- Scala
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
def makeUDF(value: Int) = udf(() => value)
var x = 123
val fooUDF = makeUDF(x) // Captures the current value
x = 456
val df = spark.range(1).select(fooUDF())
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.
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
- Scala
try:
df = ...
df.columns # This will trigger eager analysis
except Exception as e:
print(f"Error: {repr(e)}")
import org.apache.spark.SparkThrowable
import org.apache.spark.sql.functions._
val df = spark.createDataFrame(Seq(("Alice", 25), ("Bob", 30))).toDF("name", "age")
try {
val df2 = df.select("name", "age")
.withColumn("age_group", when(col("age") < 18, "minor").otherwise("adult"))
.filter(col("age_with_typo") > 6)
df2.columns // Trigger eager analysis to catch the error
} catch {
case e: SparkThrowable => println(s"Error: ${e.getMessage}")
}
É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
- Scala
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()
import org.apache.spark.sql.functions._
var df = spark.range(10).toDF
val columns = scala.collection.mutable.Set(df.columns: _*)
for (i <- 0 until 200) {
val newColumnName = i.toString
// if (!df.columns.contains(newColumnName)) { // Bad practice. The `df.columns` call causes an analysis request on the newly created DataFrame in every iteration.
if (!columns.contains(newColumnName)) { // Check the set without triggering analysis
df = df.withColumn(newColumnName, col("id") + i)
columns.add(newColumnName)
}
}
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
- Scala
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)
import org.apache.spark.sql.types.StructType
df = ...
val structColumnFields = df.schema.fields
.filter(_.dataType.isInstanceOf[StructType])
.map { field =>
// field.name -> df.select(field.name + ".*").columns // Bad practice. This creates an intermediate DataFrame and triggers analysis for each StructType column.
field.name -> field.dataType.asInstanceOf[StructType].fields.map(_.name) // Access StructType fields directly from the schema, avoiding analysis on intermediate DataFrames.
}
.toMap
println(structColumnFields)