SparkR
SparkR (R on Spark)
SparkR는 R에서 Apache Spark를 사용할 수 있게 해주는 경량 프론트엔드 R 패키지예요. R의 데이터 프레임처럼 선택, 필터링, 집계 같은 연산을 하지만, 대규모 데이터셋 위에서 분산으로 수행할 수 있도록 해줍니다. MLlib을 사용해 분산 머신러닝까지 지원해요. 다만 SparkR은 Apache Spark 4.0.0부터 더 이상 사용되지 않으며(deprecated), 향후 버전에서 제거될 예정이라는 점을 참고하세요.
출처: 문서
본문
목차 (Table of Contents)
- 개요 (Overview)
- SparkDataFrame
- 머신러닝 (Machine Learning)
- R과 Spark 사이의 데이터 타입 매핑
- 구조적 스트리밍 (Structured Streaming)
- SparkR에서의 Apache Arrow
- R 함수 이름 충돌 (R Function Name Conflicts)
- 마이그레이션 가이드 (Migration Guide)
SparkR은 Apache Spark 4.0.0부터 deprecated이며 향후 버전에서 제거될 예정입니다.
개요 (Overview)
SparkR은 R에서 Apache Spark를 사용하기 위한 경량 프론트엔드를 제공하는 R 패키지예요. Spark 4.2.0에서 SparkR은 R 데이터 프레임, dplyr과 비슷한 선택, 필터링, 집계 등의 연산을 지원하지만 대규모 데이터셋에서 동작하는 분산 데이터 프레임 구현을 제공합니다. SparkR은 MLlib을 사용한 분산 머신러닝도 지원해요.
SparkDataFrame
SparkDataFrame은 이름이 있는 열로 구성된 데이터의 분산 컬렉션이에요. 개념적으로 관계형 데이터베이스의 테이블이나 R의 데이터 프레임과 동등하지만, 내부적으로 더 풍부한 최적화를 가져요. SparkDataFrame은 구조화된 데이터 파일, Hive의 테이블, 외부 데이터베이스, 기존 로컬 R 데이터 프레임 등 다양한 소스에서 만들 수 있습니다.
이 페이지의 모든 예제는 R 또는 Spark 배포판에 포함된 샘플 데이터를 사용하며 ./bin/sparkR 셸로 실행할 수 있어요.
시작하기: SparkSession
R:
SparkR의 진입점은 R 프로그램을 Spark 클러스터에 연결하는 `SparkSession`이에요. `sparkR.session`을 사용해 `SparkSession`을 만들고 애플리케이션 이름, 의존하는 spark 패키지 등의 옵션을 전달할 수 있습니다. 또한 `SparkSession`을 통해 SparkDataFrame으로 작업할 수도 있어요. `sparkR` 셸에서 작업한다면 `SparkSession`이 이미 만들어져 있어서 `sparkR.session`을 호출할 필요가 없습니다.
sparkR.session()
RStudio에서 시작하기
RStudio에서도 SparkR을 시작할 수 있어요. RStudio, R 셸, Rscript 또는 다른 R IDE에서 R 프로그램을 Spark 클러스터에 연결할 수 있습니다. 시작하려면 환경에 SPARK_HOME이 설정되어 있는지 확인하고(Sys.getenv로 확인), SparkR 패키지를 로드한 다음 아래처럼 sparkR.session을 호출하세요. 이는 Spark 설치를 확인하고, 없으면 자동으로 다운로드·캐시합니다. 또는 install.spark를 직접 실행할 수도 있어요.
sparkR.session 호출 외에도 특정 Spark driver 속성을 지정할 수 있어요. 보통 이러한 애플리케이션 속성과 런타임 환경은 driver JVM 프로세스가 이미 시작되었으므로 프로그래밍 방식으로 설정할 수 없어요. 이 경우 SparkR이 이를 대신 처리합니다. 설정하려면 sparkR.session()의 sparkConfig 인자에 다른 설정 속성처럼 전달하면 됩니다.
if (nchar(Sys.getenv("SPARK_HOME")) < 1) {
Sys.setenv(SPARK_HOME = "/home/spark")
}
library(SparkR, lib.loc = c(file.path(Sys.getenv("SPARK_HOME"), "R", "lib")))
sparkR.session(master = "local[*]", sparkConfig = list(spark.driver.memory = "2g"))
RStudio에서 sparkR.session으로 sparkConfig에 설정할 수 있는 Spark driver 속성은 다음과 같아요:
| 속성 이름 | 속성 그룹 | spark-submit 등가물 |
|---|---|---|
spark.master |
Application Properties | --master |
spark.kerberos.keytab |
Application Properties | --keytab |
spark.kerberos.principal |
Application Properties | --principal |
spark.driver.memory |
Application Properties | --driver-memory |
spark.driver.extraClassPath |
Runtime Environment | --driver-class-path |
spark.driver.extraJavaOptions |
Runtime Environment | --driver-java-options |
spark.driver.extraLibraryPath |
Runtime Environment | --driver-library-path |
SparkDataFrames 만들기
SparkSession을 사용하면 애플리케이션이 로컬 R 데이터 프레임, Hive 테이블, 또는 다른 데이터 소스에서 SparkDataFrame을 만들 수 있어요.
로컬 데이터 프레임으로부터
데이터 프레임을 만드는 가장 간단한 방법은 로컬 R 데이터 프레임을 SparkDataFrame으로 변환하는 것이에요. 구체적으로 as.DataFrame이나 createDataFrame을 사용해 로컬 R 데이터 프레임을 전달하면 SparkDataFrame을 만들 수 있습니다. 예를 들어 다음은 R의 faithful 데이터셋을 기반으로 SparkDataFrame을 만들어요.
df <- as.DataFrame(faithful)
# Displays the first part of the SparkDataFrame
head(df)
## eruptions waiting
##1 3.600 79
##2 1.800 54
##3 3.333 74
데이터 소스로부터
SparkR은 SparkDataFrame 인터페이스를 통해 다양한 데이터 소스를 다룰 수 있어요. 이 섹션은 데이터 소스를 사용해 데이터를 로드하고 저장하는 일반적인 방법을 설명합니다. 내장 데이터 소스에 사용 가능한 더 구체적인 옵션은 Spark SQL 프로그래밍 가이드를 확인하세요.
데이터 소스에서 SparkDataFrame을 만드는 일반적인 방법은 read.df입니다. 이 메서드는 로드할 파일의 경로와 데이터 소스 유형을 받으며, 현재 활성화된 SparkSession을 자동으로 사용해요. SparkR은 JSON, CSV, Parquet 파일을 기본적으로 읽으며, Third Party Projects 같은 소스에서 제공되는 패키지를 통해 Avro 같은 인기 있는 파일 형식의 데이터 소스 커넥터를 찾을 수 있어요. 이러한 패키지는 spark-submit이나 sparkR 명령에 --packages를 지정해 추가하거나, 대화형 R 셸이나 RStudio에서 sparkPackages 파라미터로 SparkSession을 초기화할 때 지정할 수 있습니다.
sparkR.session(sparkPackages = "org.apache.spark:spark-avro_2.13:4.2.0")
예제 JSON 입력 파일을 사용해 데이터 소스를 어떻게 사용하는지 볼 수 있어요. 여기서 사용되는 파일은 전형적인 JSON 파일이 아니라는 점을 주의하세요. 파일의 각 줄은 별개의, 자체 포함된 유효한 JSON 객체를 포함해야 합니다. 자세한 내용은 JSON Lines 텍스트 형식, newline-delimited JSON을 참고하세요. 결과적으로 일반적인 여러 줄 JSON 파일은 대부분 실패합니다.
people <- read.df("./examples/src/main/resources/people.json", "json")
head(people)
## age name
##1 NA Michael
##2 30 Andy
##3 19 Justin
# SparkR automatically infers the schema from the JSON file
printSchema(people)
# root
# |-- age: long (nullable = true)
# |-- name: string (nullable = true)
# Similarly, multiple files can be read with read.json
people <- read.json(c("./examples/src/main/resources/people.json", "./examples/src/main/resources/people2.json"))
데이터 소스 API는 CSV 형식의 입력 파일을 기본적으로 지원해요. 자세한 내용은 SparkR read.df API 문서를 참고하세요.
df <- read.df(csvPath, "csv", header = "true", inferSchema = "true", na.strings = "NA")
데이터 소스 API는 SparkDataFrame을 여러 파일 형식으로 저장하는 데도 사용할 수 있어요. 예를 들어 이전 예제의 SparkDataFrame을 write.df를 사용해 Parquet 파일로 저장할 수 있습니다.
write.df(people, path = "people.parquet", source = "parquet", mode = "overwrite")
Hive 테이블로부터
Hive 테이블에서도 SparkDataFrame을 만들 수 있어요. 그러려면 Hive MetaStore의 테이블에 접근할 수 있는 Hive 지원 SparkSession을 만들어야 합니다. 참고로 Spark는 Hive 지원으로 빌드되어야 하며, 자세한 내용은 SQL 프로그래밍 가이드에서 확인할 수 있어요. SparkR에서 기본적으로 Hive 지원을 활성화한(enableHiveSupport = TRUE) SparkSession을 만들려고 시도합니다.
sparkR.session()
sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING)")
sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src")
# Queries can be expressed in HiveQL.
results <- sql("FROM src SELECT key, value")
# results is now a SparkDataFrame
head(results)
## key value
## 1 238 val_238
## 2 86 val_86
## 3 311 val_311
SparkDataFrame 연산
SparkDataFrame은 구조화된 데이터 처리를 위해 많은 함수를 지원해요. 여기에 기본적인 예제 몇 가지를 포함하고, 전체 목록은 API 문서에서 찾을 수 있습니다.
행, 열 선택
# Create the SparkDataFrame
df <- as.DataFrame(faithful)
# Get basic information about the SparkDataFrame
df
## SparkDataFrame[eruptions:double, waiting:double]
# Select only the "eruptions" column
head(select(df, df$eruptions))
## eruptions
##1 3.600
##2 1.800
##3 3.333
# You can also pass in column name as strings
head(select(df, "eruptions"))
# Filter the SparkDataFrame to only retain rows with wait times shorter than 50 mins
head(filter(df, df$waiting < 50))
## eruptions waiting
##1 1.750 47
##2 1.750 47
##3 1.867 48
그룹화, 집계
SparkR 데이터 프레임은 그룹화 후 데이터를 집계하는 데 흔히 사용되는 여러 함수를 지원해요. 예를 들어 faithful 데이터셋의 waiting 시간 히스토그램을 아래처럼 계산할 수 있습니다.
# We use the `n` operator to count the number of times each waiting time appears
head(summarize(groupBy(df, df$waiting), count = n(df$waiting)))
## waiting count
##1 70 4
##2 67 1
##3 69 2
# We can also sort the output from the aggregation to get the most common waiting times
waiting_counts <- summarize(groupBy(df, df$waiting), count = n(df$waiting))
head(arrange(waiting_counts, desc(waiting_counts$count)))
## waiting count
##1 78 15
##2 83 14
##3 81 13
표준 집계 외에도 SparkR은 OLAP 큐브 연산자 cube를 지원해요:
head(agg(cube(df, "cyl", "disp", "gear"), avg(df$mpg)))
## cyl disp gear avg(mpg)
##1 NA 140.8 4 22.8
##2 4 75.7 4 30.4
##3 8 400.0 3 19.2
##4 8 318.0 3 15.5
##5 NA 351.0 NA 15.8
##6 NA 275.8 NA 16.3
그리고 rollup:
head(agg(rollup(df, "cyl", "disp", "gear"), avg(df$mpg)))
## cyl disp gear avg(mpg)
##1 4 75.7 4 22.8
##2 8 400.0 3 19.2
##3 8 318.0 3 15.5
##4 4 78.7 NA 32.4
##5 8 304.0 3 15.2
##6 4 79.0 NA 27.3
열 연산
SparkR은 데이터 처리와 집계 중에 열에 직접 적용할 수 있는 여러 함수도 제공해요. 아래 예제는 기본 산술 함수의 사용을 보여줍니다.
# Convert waiting time from hours to seconds.
# Note that we can assign this to a new column in the same SparkDataFrame
df$waiting_secs <- df$waiting * 60
head(df)
## eruptions waiting waiting_secs
##1 3.600 79 4740
##2 1.800 54 3240
##3 3.333 74 4440
사용자 정의 함수 적용
SparkR에서는 여러 종류의 사용자 정의 함수를 지원해요.
dapply 또는 dapplyCollect로 대규모 데이터셋에 함수 실행
dapply — SparkDataFrame의 각 파티션에 함수를 적용합니다. SparkDataFrame의 각 파티션에 적용될 함수는 파라미터가 하나뿐이어야 하며, 각 파티션에 해당하는 data.frame이 전달됩니다. 함수의 출력은 data.frame이어야 해요. Schema는 결과 SparkDataFrame의 행 형식을 지정합니다. 반환 값의 데이터 타입과 일치해야 합니다.
# Convert waiting time from hours to seconds.
# Note that we can apply UDF to DataFrame.
schema <- structType(structField("eruptions", "double"), structField("waiting", "double"),
structField("waiting_secs", "double"))
df1 <- dapply(df, function(x) { x <- cbind(x, x$waiting * 60) }, schema)
head(collect(df1))
## eruptions waiting waiting_secs
##1 3.600 79 4740
##2 1.800 54 3240
##3 3.333 74 4440
##4 2.283 62 3720
##5 4.533 85 5100
##6 2.883 55 3300
dapplyCollect — dapply처럼 SparkDataFrame의 각 파티션에 함수를 적용하고 결과를 다시 수집합니다. 함수의 출력은 data.frame이어야 해요. 하지만 schema는 전달할 필요가 없습니다. 참고로 dapplyCollect는 모든 파티션에서 실행된 UDF의 출력을 드라이버로 가져와 driver 메모리에 맞출 수 없으면 실패할 수 있어요.
# Convert waiting time from hours to seconds.
# Note that we can apply UDF to DataFrame and return a R's data.frame
ldf <- dapplyCollect(
df,
function(x) {
x <- cbind(x, "waiting_secs" = x$waiting * 60)
})
head(ldf, 3)
## eruptions waiting waiting_secs
##1 3.600 79 4740
##2 1.800 54 3240
##3 3.333 74 4440
입력 열(들)로 그룹화하고 gapply 또는 gapplyCollect로 함수 실행
gapply — SparkDataFrame의 각 그룹에 함수를 적용합니다. SparkDataFrame의 각 그룹에 적용될 함수는 파라미터가 정확히 두 개여야 해요: 그룹 키와 그 키에 해당하는 R data.frame. 그룹은 SparkDataFrame의 열(들)에서 선택됩니다. 함수의 출력은 data.frame이어야 해요. Schema는 결과 SparkDataFrame의 행 형식을 지정합니다. Spark 데이터 타입을 기준으로 R 함수의 출력 schema를 나타내야 합니다. 반환된 data.frame의 열 이름은 사용자가 설정합니다.
# Determine six waiting times with the largest eruption time in minutes.
schema <- structType(structField("waiting", "double"), structField("max_eruption", "double"))
result <- gapply(
df,
"waiting",
function(key, x) {
y <- data.frame(key, max(x$eruptions))
},
schema)
head(collect(arrange(result, "max_eruption", decreasing = TRUE)))
## waiting max_eruption
##1 64 5.100
##2 69 5.067
##3 71 5.033
##4 87 5.000
##5 63 4.933
##6 89 4.900
gapplyCollect — gapply처럼 SparkDataFrame의 각 파티션에 함수를 적용하고 결과를 R data.frame으로 다시 수집합니다. 함수의 출력은 data.frame이어야 해요. 하지만 schema는 전달할 필요가 없습니다. 참고로 gapplyCollect는 모든 파티션에서 실행된 UDF의 출력을 드라이버로 가져와 driver 메모리에 맞출 수 없으면 실패할 수 있어요.
# Determine six waiting times with the largest eruption time in minutes.
result <- gapplyCollect(
df,
"waiting",
function(key, x) {
y <- data.frame(key, max(x$eruptions))
colnames(y) <- c("waiting", "max_eruption")
y
})
head(result[order(result$max_eruption, decreasing = TRUE), ])
## waiting max_eruption
##1 64 5.100
##2 69 5.067
##3 71 5.033
##4 87 5.000
##5 63 4.933
##6 89 4.900
spark.lapply로 로컬 R 함수 분산 실행
spark.lapply — 네이티브 R의 lapply와 비슷하게, spark.lapply는 요소 목록에 대해 함수를 실행하고 Spark로 계산을 분산합니다. doParallel이나 lapply와 비슷한 방식으로 목록의 요소에 함수를 적용해요. 모든 계산의 결과는 단일 머신에 맞아야 합니다. 그렇지 않으면 df <- createDataFrame(list) 같은 것을 한 다음 dapply를 사용할 수 있어요.
# Perform distributed training of multiple models with spark.lapply. Here, we pass
# a read-only list of arguments which specifies family the generalized linear model should be.
families <- c("gaussian", "poisson")
train <- function(family) {
model <- glm(Sepal.Length ~ Sepal.Width + Species, iris, family = family)
summary(model)
}
# Return a list of model's summaries
model.summaries <- spark.lapply(families, train)
# Print the summary of each model
print(model.summaries)
즉시 실행 (Eager execution)
즉시 실행이 활성화되면 SparkDataFrame이 만들어질 때 데이터가 즉시 R 클라이언트로 반환돼요. 기본적으로 즉시 실행은 활성화되지 않으며, SparkSession이 시작될 때 설정 속성 spark.sql.repl.eagerEval.enabled를 true로 설정해 활성화할 수 있습니다.
표시할 데이터의 최대 행 수와 열당 최대 문자 수는 각각 설정 속성 spark.sql.repl.eagerEval.maxNumRows와 spark.sql.repl.eagerEval.truncate로 제어할 수 있어요. 이 속성들은 즉시 실행이 활성화된 경우에만 유효합니다. 명시적으로 설정하지 않으면 기본적으로 최대 20행, 열당 최대 20문자가 표시됩니다.
# Start up spark session with eager execution enabled
sparkR.session(master = "local[*]",
sparkConfig = list(spark.sql.repl.eagerEval.enabled = "true",
spark.sql.repl.eagerEval.maxNumRows = as.integer(10)))
# Create a grouped and sorted SparkDataFrame
df <- createDataFrame(faithful)
df2 <- arrange(summarize(groupBy(df, df$waiting), count = n(df$waiting)), "waiting")
# Similar to R data.frame, displays the data returned, instead of SparkDataFrame class string
df2
##+-------+-----+
##|waiting|count|
##+-------+-----+
##| 43.0| 1|
##| 45.0| 3|
##| 46.0| 5|
##| 47.0| 4|
##| 48.0| 3|
##| 49.0| 5|
##| 50.0| 5|
##| 51.0| 6|
##| 52.0| 5|
##| 53.0| 7|
##+-------+-----+
##only showing top 10 rows
sparkR 셸에서 즉시 실행을 활성화하려면 --conf 옵션에 spark.sql.repl.eagerEval.enabled=true 설정 속성을 추가하세요.
SparkR에서 SQL 쿼리 실행
SparkDataFrame은 Spark SQL에서 임시 뷰로 등록될 수 있고, 그러면 그 데이터에 대해 SQL 쿼리를 실행할 수 있어요. sql 함수는 애플리케이션이 프로그래밍 방식으로 SQL 쿼리를 실행하고 결과를 SparkDataFrame으로 반환할 수 있게 해줍니다.
# Load a JSON file
people <- read.df("./examples/src/main/resources/people.json", "json")
# Register this SparkDataFrame as a temporary view.
createOrReplaceTempView(people, "people")
# SQL statements can be run by using the sql method
teenagers <- sql("SELECT name FROM people WHERE age >= 13 AND age <= 19")
head(teenagers)
## name
##1 Justin
머신러닝 (Machine Learning)
알고리즘 (Algorithms)
SparkR은 현재 다음 머신러닝 알고리즘을 지원해요.
분류 (Classification)
spark.logit:로지스틱 회귀 (Logistic Regression)spark.mlp:다층 퍼셉트론 (Multilayer Perceptron, MLP)spark.naiveBayes:나이브 베이즈 (Naive Bayes)spark.svmLinear:선형 서포트 벡터 머신 (Linear Support Vector Machine)spark.fmClassifier:분해 머신 분류기 (Factorization Machines classifier)
회귀 (Regression)
spark.survreg:가속 실패 시간 (AFT) 생존 모델 (Accelerated Failure Time Survival Model)spark.glm또는glm:일반화 선형 모델 (Generalized Linear Model, GLM)spark.isoreg:등장 회귀 (Isotonic Regression)spark.lm:선형 회귀 (Linear Regression)spark.fmRegressor:분해 머신 회귀 (Factorization Machines regressor)
트리 (Tree)
spark.decisionTree:회귀및분류용Decision Treespark.gbt:회귀및분류용Gradient Boosted Treesspark.randomForest:회귀및분류용Random Forest
클러스터링 (Clustering)
spark.bisectingKmeans:이분 k-means (Bisecting k-means)spark.gaussianMixture:가우시안 혼합 모델 (Gaussian Mixture Model, GMM)spark.kmeans:K-Meansspark.lda:잠재 디리클레 할당 (Latent Dirichlet Allocation, LDA)spark.powerIterationClustering (PIC):전력 반복 클러스터링 (Power Iteration Clustering, PIC)
협업 필터링 (Collaborative Filtering)
빈발 패턴 마이닝 (Frequent Pattern Mining)
통계 (Statistics)
spark.kstest:Kolmogorov-Smirnov Test
내부적으로 SparkR은 MLlib을 사용해 모델을 훈련합니다. 예제 코드는 MLlib 사용자 가이드의 해당 섹션을 참고하세요. 사용자는 summary를 호출해 피팅된 모델 요약을 출력하고, predict로 새 데이터를 예측하며, write.ml/read.ml로 피팅된 모델을 저장/로드할 수 있어요. SparkR은 모델 피팅에 사용 가능한 R 공식 연산자의 일부('~', '.', ':', '+', '-')를 지원합니다.
모델 영속화 (Model persistence)
다음 예제는 SparkR로 MLlib 모델을 저장/로드하는 방법을 보여줘요.
training <- read.df("data/mllib/sample_multiclass_classification_data.txt", source = "libsvm")
# Fit a generalized linear model of family "gaussian" with spark.glm
df_list <- randomSplit(training, c(7,3), 2)
gaussianDF <- df_list[[1]]
gaussianTestDF <- df_list[[2]]
gaussianGLM <- spark.glm(gaussianDF, label ~ features, family = "gaussian")
# Save and then load a fitted MLlib model
modelPath <- tempfile(pattern = "ml", fileext = ".tmp")
write.ml(gaussianGLM, modelPath)
gaussianGLM2 <- read.ml(modelPath)
# Check model summary
summary(gaussianGLM2)
# Check model prediction
gaussianPredictions <- predict(gaussianGLM2, gaussianTestDF)
head(gaussianPredictions)
unlink(modelPath)
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/ml.R"에서 확인할 수 있어요.
R과 Spark 사이의 데이터 타입 매핑
| R | Spark |
|---|---|
| byte | byte |
| integer | integer |
| float | float |
| double | double |
| numeric | double |
| character | string |
| string | string |
| binary | binary |
| raw | binary |
| logical | boolean |
| POSIXct | timestamp |
| POSIXlt | timestamp |
| Date | date |
| array | array |
| list | array |
| env | map |
구조적 스트리밍 (Structured Streaming)
SparkR은 Structured Streaming API를 지원해요. Structured Streaming은 Spark SQL 엔진 위에 구축된 확장 가능하고 장애에 강한 스트림 처리 엔진입니다. 자세한 내용은 Structured Streaming 프로그래밍 가이드의 R API를 참고하세요.
SparkR에서의 Apache Arrow
Apache Arrow는 JVM과 R 프로세스 간에 데이터를 효율적으로 전송하기 위해 Spark에서 사용하는 인-메모리 컬럼 지향 데이터 형식이에요. PySpark 최적화에 대해서도 참고하세요: PySpark Usage Guide for Pandas with Apache Arrow. 이 가이드는 SparkR에서 Arrow 최적화를 몇 가지 핵심 포인트와 함께 어떻게 사용하는지 설명하는 것을 목표로 해요.
Arrow 설치 확인 (Ensure Arrow Installed)
Arrow R 라이브러리는 CRAN에서 사용 가능하며 아래처럼 설치할 수 있어요.
Rscript -e 'install.packages("arrow", repos="https://cloud.r-project.org/")'
자세한 내용은 Apache Arrow 공식 문서를 참고하세요.
Arrow R 패키지가 모든 클러스터 노드에 설치되어 사용 가능해야 한다는 점을 반드시 확인하세요. 현재 지원되는 최소 버전은 1.0.0이지만, SparkR의 Arrow 최적화는 실험적이므로 마이너 릴리스 사이에 바뀔 수 있어요.
R DataFrame과의 변환, dapply, gapply 활성화
Arrow 최적화는 collect(spark_df) 호출로 Spark DataFrame을 R DataFrame으로 변환할 때, createDataFrame(r_df)로 R DataFrame에서 Spark DataFrame을 만들 때, dapply(...)로 각 파티션에 R 네이티브 함수를 적용할 때, gapply(...)로 그룹화된 데이터에 R 네이티브 함수를 적용할 때 사용할 수 있어요. 이를 실행할 때 Arrow를 사용하려면 먼저 Spark 설정 'spark.sql.execution.arrow.sparkr.enabled'를 'true'로 설정해야 해요. 기본적으로 활성화되어 있지 않습니다.
최적화가 활성화되어 있든 아니든 SparkR은 동일한 결과를 만들어요. 또한 Spark DataFrame과 R DataFrame 사이의 변환은 실제 계산 전에 최적화가 어떤 이유로든 실패하면 자동으로 비-Arrow 최적화 구현으로 폴백합니다.
# Start up spark session with Arrow optimization enabled
sparkR.session(master = "local[*]",
sparkConfig = list(spark.sql.execution.arrow.sparkr.enabled = "true"))
# Converts Spark DataFrame from an R DataFrame
spark_df <- createDataFrame(mtcars)
# Converts Spark DataFrame to an R DataFrame
collect(spark_df)
# Apply an R native function to each partition.
collect(dapply(spark_df, function(rdf) { data.frame(rdf$gear + 1) }, structType("gear double")))
# Apply an R native function to grouped data.
collect(gapply(spark_df,
"gear",
function(key, group) {
data.frame(gear = key[[1]], disp = mean(group$disp) > group$disp)
},
structType("gear double, disp boolean")))
참고로 Arrow를 사용해도 collect(spark_df)는 DataFrame의 모든 레코드를 드라이버 프로그램으로 수집하므로 데이터의 작은 부분집합에 대해 수행해야 합니다. 또한 gapply(...)와 dapply(...)에서 지정된 출력 schema는 주어진 함수가 반환하는 R DataFrame과 일치해야 해요.
지원되는 SQL 타입 (Supported SQL Types)
현재 Arrow 기반 변환은 FloatType, BinaryType, ArrayType, StructType, MapType을 제외한 모든 Spark SQL 데이터 타입을 지원해요.
R 함수 이름 충돌 (R Function Name Conflicts)
R에서 새 패키지를 로드하고 부착할 때, 함수가 다른 함수를 가리는(masking) 이름 충돌이 발생할 수 있어요.
다음 함수들은 SparkR 패키지에 의해 가려집니다:
| 가려진 함수 | 접근 방법 |
|---|---|
cov in package:stats |
stats::cov(x, y = NULL, use = "everything", method = c("pearson", "kendall", "spearman")) |
filter in package:stats |
stats::filter(x, filter, method = c("convolution", "recursive"), sides = 2, circular = FALSE, init) |
sample in package:base |
base::sample(x, size, replace = FALSE, prob = NULL) |
SparkR의 일부는 dplyr 패키지를 모델로 했기 때문에, SparkR의 특정 함수는 dplyr의 함수와 같은 이름을 공유해요. 두 패키지의 로드 순서에 따라, 먼저 로드된 패키지의 일부 함수가 나중에 로드된 패키지의 함수에 의해 가려집니다. 이런 경우 패키지 이름을 접두사로 붙여 호출하세요. 예를 들어 SparkR::cume_dist(x) 또는 dplyr::cume_dist(x)처럼요.
R의 검색 경로는 search()로 확인할 수 있어요.
마이그레이션 가이드 (Migration Guide)
마이그레이션 가이드는 이제 이 페이지에 보관되어 있어요.
더 알아보기 (Learn more)
- SparkR 마이그레이션 가이드: SparkR의 버전별 변경 사항을 확인해요.
- MLlib 주요 가이드 (MLlib: Main Guide): SparkR이 사용하는 분산 머신러닝 구조를 이해해요.
- Spark SQL 프로그래밍 가이드: DataFrame과 SQL의 기초를 함께 익혀보세요.