Aller au contenu principal

Migrer de SparkR vers sparklyr

SparkR a été développé dans le cadre d'Apache Spark, et sa conception est familière aux utilisateurs de Scala et Python, mais potentiellement moins intuitive pour les praticiens de R. De plus, SparkR est obsolète dans Spark 4,0.

En revanche, sparklyr vise à offrir une expérience plus conviviale avec R. Il tire parti de la syntaxe dplyr, familière aux utilisateurs de tidyverse avec des modèles comme select(), filter() et mutate() pour les Opérations de DataFrame.

sparklyr est le paquet R recommandé pour travailler avec Apache Spark. Cette page explique les différences entre SparkR et sparklyr au sein des Spark API, et fournit des informations sur la migration de code.

Configuration de l'environnement

Installation

Si vous êtes dans le Workspace Databricks, aucune installation n'est requise. Chargez sparklyr avec library(sparklyr). Pour installer sparklyr localement en dehors de Databricks, consultez Démarrer.

Connexion à Spark

Connectez-vous à Spark avec sparklyr dans le workspace Databricks ou localement à l'aide de Databricks Connect:

Workspace :

R
library(sparklyr)
sc <- spark_connect(method = "databricks")

Databricks Connect :

R
sc <- spark_connect(method = "databricks_connect")

Pour plus de détails et un tutoriel étendu sur Databricks Connect avec sparklyr, consultez Premiers pas.

Lecture et écriture des données

sparklyr dispose d'une famille de fonctions spark_read_*() et spark_write_*() pour charger et enregistrer des données, contrairement aux fonctions génériques read.df() et write.df() de SparkR. Il existe également des fonctions uniques pour créer des DataFrames Spark ou des vues temporaires Spark SQL à partir de DataFrames R en mémoire.

Tâche

sparkr

sparklyr

Copier les données vers Spark

createDataFrame()

copy_to()

Créer une vue temporaire

createOrReplaceTempView()

Utilisez invoke() directement avec la méthode

Écrire des données dans une table

saveAsTable()

spark_write_table()

Écrire des données dans un format spécifié.

write.df()

spark_write_<format>()

Lire les données de la table

tableToDF()

tbl() OU spark_read_table()

Lire les données à partir d'un format spécifié

read.df()

spark_read_<format>()

Tâche

sparkr

sparklyr

Copier les données vers Spark

createDataFrame()

copy_to()

Créer une vue temporaire

createOrReplaceTempView()

Utilisez invoke() directement avec la méthode

Écrire des données dans une table

saveAsTable()

spark_write_table()

Écrire des données dans un format spécifié.

write.df()

spark_write_<format>()

Lire les données de la table

tableToDF()

tbl() OU spark_read_table()

Lire les données à partir d'un format spécifié

read.df()

spark_read_<format>()

Chargement des données

Pour convertir un cadre de données R en un Spark DataFrame, ou pour créer une vue temporaire à partir d'un DataFrame afin de lui appliquer du SQL :

R
mtcars_df <- createDataFrame(mtcars)

copy_to() crée une vue temporaire en utilisant le nom spécifié. Vous pouvez utiliser le nom pour référencer les données si vous utilisez SQL directement (par exemple, sdf_sql()). De plus, copy_to() met en cache les données en définissant le paramètre memory sur TRUE.

Création de vues

Les exemples de code suivants montrent comment les vues temporaires sont créées :

R
createOrReplaceTempView(mtcars_df, "mtcars_tmp_view")

Écriture de données

Les exemples de code suivants montrent comment les données sont écrites :

R
# Save a DataFrame to Unity Catalog
saveAsTable(
mtcars_df,
tableName = "<catalog>.<schema>.<table>",
mode = "overwrite"
)

# Save a DataFrame to local filesystem using Delta format
write.df(
mtcars_df,
path = "file:/<path/to/save/delta/mtcars>",
source = "delta",
mode = "overwrite"
)

Lecture des données

Les exemples de code suivants montrent comment les données sont lues :

R
# Load a Unity Catalog table as a DataFrame
tableToDF("<catalog>.<schema>.<table>")

# Load csv file into a DataFrame
read.df(
path = "file:/<path/to/read/csv/data.csv>",
source = "csv",
header = TRUE,
inferSchema = TRUE
)

# Load Delta from local filesystem as a DataFrame
read.df(
path = "file:/<path/to/read/delta/mtcars>",
source = "delta"
)

# Load data from a table using SQL - Databricks recommendeds using `tableToDF`
sql("SELECT * FROM <catalog>.<schema>.<table>")

Traitement des données

Sélectionner et filtrer

R
# Select specific columns
select(mtcars_df, "mpg", "cyl", "hp")

# Filter rows where mpg > 20
filter(mtcars_df, mtcars_df$mpg > 20)

Ajouter des colonnes

R
# Add a new column 'power_to_weight' (hp divided by wt)
withColumn(mtcars_df, "power_to_weight", mtcars_df$hp / mtcars_df$wt)

Regroupement et agrégation

R
# Calculate average mpg and hp by number of cylinders
mtcars_df |>
groupBy("cyl") |>
summarize(
avg_mpg = avg(mtcars_df$mpg),
avg_hp = avg(mtcars_df$hp)
)

Jointures

Supposons que nous ayons un autre dataset avec des étiquettes de cylindre que nous voulons joindre à mtcars.

R
# Create another DataFrame with cylinder labels
cylinders <- data.frame(
cyl = c(4, 6, 8),
cyl_label = c("Four", "Six", "Eight")
)
cylinders_df <- createDataFrame(cylinders)

# Join mtcars_df with cylinders_df
join(
x = mtcars_df,
y = cylinders_df,
mtcars_df$cyl == cylinders_df$cyl,
joinType = "inner"
)

Fonctions définies par l'utilisateur (UDF)

Pour créer une fonction personnalisée pour la catégorisation :

R
# Define the custom function
categorize_hp <- function(df)
df$hp_category <- ifelse(df$hp > 150, "High", "Low") # a real-world example would use case_when() with mutate()
df

SparkR nécessite de définir explicitement le schéma de sortie avant d'appliquer une fonction :

R
# Define the schema for the output DataFrame
schema <- structType(
structField("mpg", "double"),
structField("cyl", "double"),
structField("disp", "double"),
structField("hp", "double"),
structField("drat", "double"),
structField("wt", "double"),
structField("qsec", "double"),
structField("vs", "double"),
structField("am", "double"),
structField("gear", "double"),
structField("carb", "double"),
structField("hp_category", "string")
)

# Apply the function across partitions
dapply(
mtcars_df,
func = categorize_hp,
schema = schema
)

# Apply the same function to each group of a DataFrame. Note that the schema is still required.
gapply(
mtcars_df,
cols = "hp",
func = categorize_hp,
schema = schema
)

spark.lapply() vs. spark_apply()

Dans SparkR, spark.lapply() fonctionne sur les listes R plutôt que sur les DataFrames. Il n'y a pas d'équivalent direct dans sparklyr, mais vous pouvez obtenir un comportement similaire avec spark_apply() en travaillant avec un DataFrame qui inclut des identifiants uniques et en regroupant par ces ID. Dans certains cas, les opérations par ligne peuvent également fournir des fonctionnalités comparables. Pour plus d'informations sur spark_apply(), voir Distribution des calculs R.

R
# Define a list of integers
numbers <- list(1, 2, 3, 4, 5)

# Define a function to apply
square <- function(x)
x * x

# Apply the function over list using Spark
spark.lapply(numbers, square)

Machine Learning

Des exemples complets de SparkR et sparklyr pour le machine learning se trouvent dans le Guide ML de Spark et la référence sparklyr.

remarque

Si vous n'utilisez pas Spark MLlib, Databricks recommande d'utiliser des UDF pour l'entraînement avec la bibliothèque de votre choix (par exemple xgboost).

Régression linéaire

R
# Select features
training_df <- select(mtcars_df, "mpg", "hp", "wt")

# Fit the model using Generalized Linear Model (GLM)
linear_model <- spark.glm(training_df, mpg ~ hp + wt, family = "gaussian")

# View model summary
summary(linear_model)

Regroupement K-moyennes

R
# Apply KMeans clustering with 3 clusters using mpg and hp as features
kmeans_model <- spark.kmeans(mtcars_df, mpg ~ hp, k = 3)

# Get cluster predictions
predict(kmeans_model, mtcars_df)

Performances et optimisation

Collecte

SparkR et sparklyr utilisent tous deux collect() pour convertir les DataFrames Spark en DataFrames R. Ne collectez que de petites quantités de données dans les DataFrames R, sinon le Driver Spark manquera de mémoire.

Pour éviter les erreurs de mémoire insuffisante, SparkR dispose d'optimisations intégrées dans Databricks Runtime qui aident à collecter des données ou à exécuter des fonctions définies par l'utilisateur.

Pour garantir des performances optimales avec sparklyr pour la collecte de données et les fonctions UDF sur les versions de Databricks Runtime antérieures à 14.3 LTS, chargez le package arrow :

R
library(arrow)

Partitionnement en mémoire

R
# Repartition the SparkDataFrame based on 'cyl' column
repartition(mtcars_df, col = mtcars_df$cyl)

# Repartition the SparkDataFrame to number of partitions
repartition(mtcars_df, numPartitions = 10)

# Coalesce the DataFrame to number of partitions
coalesce(mtcars_df, numPartitions = 1)

# Get number of partitions
getNumPartitions(mtcars_df)

Mise en cache

R
# Cache the DataFrame in memory
cache(mtcars_df)