Spark Connect 개요

Spark Connect 개요 (Spark Connect Overview)

Spark 3.4에서 도입된 Spark Connect는 클라이언트-서버가 분리된 아키텍처로, DataFrame API와 미해결 논리 계획(unresolved logical plans)을 프로토콜로 사용해 Spark 클러스터에 원격으로 연결할 수 있게 해주는 기능이에요. Spark Connect가 어떻게 동작하는지, 어떤 이점이 있는지, 그리고 PySpark·Scala 클라이언트 애플리케이션에서 어떻게 사용하는지 알아볼게요.

출처: 문서

본문

Spark 3.4에서 Spark Connect는 분리된 클라이언트-서버 아키텍처를 도입했어요. DataFrame API와 미해결 논리 계획(unresolved logical plans)을 프로토콜로 사용해 Spark 클러스터에 원격으로 연결할 수 있게 해줘요. 클라이언트와 서버의 분리는 Spark와 그 오픈 생태계를 어디서든 활용할 수 있게 해줘요. 모던 데이터 애플리케이션, IDE, 노트북, 프로그래밍 언어에 내장될 수 있어요.

시작하려면 Quickstart: Spark Connect를 참고하세요.

Spark Connect의 동작 방식 (How Spark Connect works)

Spark Connect 클라이언트 라이브러리는 Spark 애플리케이션 개발을 단순화하도록 설계됐어요. 애플리케이션 서버, IDE, 노트북, 프로그래밍 언어 어디에나 내장될 수 있는 얇은 API예요. Spark Connect API는 클라이언트와 Spark 드라이버 사이의 언어에 구애받지 않는(language-agnostic) 프로토콜로 미해결 논리 계획(unresolved logical plans)을 사용해 Spark의 DataFrame API 위에 구축돼요.

Spark Connect 클라이언트는 DataFrame 연산을 protocol buffers로 인코딩된 미해결 논리 쿼리 계획으로 변환해요. 이들은 gRPC 프레임워크를 사용해 서버로 전송돼요.

Spark 서버에 내장된 Spark Connect 엔드포인트는 미해결 논리 계획을 받아 Spark의 논리 계획 연산자로 변환해요. 이는 SQL 쿼리를 파싱하는 것과 비슷한데, 속성과 관계가 파싱되고 초기 파스 계획이 구축돼요. 거기서부터 표준 Spark 실행 프로세스가 시작되어, Spark Connect가 Spark의 모든 최적화와 개선을 활용하도록 보장해요. 결과는 Apache Arrow로 인코딩된 행 배치들로 gRPC를 통해 클라이언트로 스트리밍돼요.

Spark Connect 클라이언트 애플리케이션이 기존 Spark 애플리케이션과 다른 점 (How Spark Connect client applications differ from classic Spark applications)

Spark Connect의 주요 설계 목표 중 하나는 클라이언트와 서버의 완전한 분리·격리(isolation)를 가능하게 하는 것이에요. 그 결과, Spark Connect를 사용할 때 개발자가 알아야 할 몇 가지 변경 사항이 있어요:

  • 클라이언트는 Spark 드라이버와 같은 프로세스에서 실행되지 않아요. 즉 클라이언트는 실행 환경을 조작하기 위해 드라이버 JVM에 직접 접근·상호작용할 수 없어요. 특히 PySpark에서 클라이언트는 Py4J를 사용하지 않으므로, DataFrame, Column, SparkSession 등의 JVM 구현을 담고 있는 private 필드(예: df._jdf)에 접근할 수 없어요.
  • 설계상 Spark Connect 프로토콜은 Spark의 논리 계획을 서버에서 실행할 연산을 선언적으로 설명하는 추상화로 사용해요. 따라서 Spark Connect 프로토콜은 Spark의 모든 실행 API를 지원하지 않는데, 가장 중요한 것은 RDD예요.
  • Spark Connect는 사용자에게 세션 기반(session-based) 클라이언트를 제공해요. 즉 클라이언트는 연결된 모든 클라이언트의 환경을 조작하는 클러스터의 속성에 접근할 수 없어요. 가장 중요하게, 클라이언트는 정적 Spark 구성이나 SparkContext에 접근할 수 없어요.

Spark Connect의 운영 이점 (Operational benefits of Spark Connect)

이 새 아키텍처로 Spark Connect는 여러 멀티테넌트 운영 문제를 완화해요:

  • 안정성 (Stability): 메모리를 너무 많이 사용하는 애플리케이션은 이제 자신만의 프로세스에서 실행될 수 있으므로 자신의 환경에만 영향을 미쳐요. 사용자는 클라이언트에서 자신만의 의존성을 정의할 수 있고, Spark 드라이버와의 잠재적 충돌을 걱정할 필요가 없어요.
  • 업그레이드 가능성 (Upgradability): Spark 드라이버는 이제 애플리케이션과 독립적으로 매끄럽게 업그레이드될 수 있어요. 예를 들어 성능 개선과 보안 패치의 이점을 얻기 위함이에요. 즉 서버 측 RPC 정의가 하위 호환되도록 설계되어 있다면, 애플리케이션은 전방 호환(forward-compatible)될 수 있어요.
  • 디버깅·관찰 가능성 (Debuggability and observability): Spark Connect는 개발 중에 선호하는 IDE에서 직접 대화형 디버깅을 가능하게 해요. 마찬가지로 애플리케이션은 애플리케이션 프레임워크의 네이티브 메트릭과 로깅 라이브러리를 사용해 모니터링할 수 있어요.

Spark Connect 사용 방법 (How to use Spark Connect)

Spark Connect는 PySpark와 Scala 애플리케이션을 지원해요. Spark Connect로 Apache Spark 서버를 실행하고, Spark Connect 클라이언트 라이브러리를 사용하는 클라이언트 애플리케이션에서 여기에 연결하는 방법을 차근차근 알아볼게요.

Spark 서버 다운로드 및 Spark Connect로 시작하기 (Download and start Spark server with Spark Connect)

먼저 Download Apache Spark 페이지에서 Spark를 다운로드하세요. 페이지 상단의 릴리스 드롭다운에서 최신 릴리스를 선택하세요. 그다음 패키지 타입, 주로 "Pre-built for Apache Hadoop 3.5 and later"를 선택하고 다운로드 링크를 클릭하세요.

이제 방금 다운로드한 Spark 패키지를 컴퓨터에서 압축 해제하세요. 예를 들어:

tar -xvf spark-4.2.0-bin-hadoop3.tgz

터미널 창에서 Spark를 압축 해제한 위치의 spark 폴더로 이동한 후 start-connect-server.sh 스크립트를 실행해 Spark Connect로 Spark 서버를 시작하세요. 예를 들면:

./sbin/start-connect-server.sh

이전에 다운로드한 Spark 버전과 같은 버전의 패키지를 사용하는지 확인하세요. 이 예제에서는 Scala 2.13이 포함된 Spark 4.2.0을 사용했어요.

이제 Spark 서버가 실행 중이며 클라이언트 애플리케이션에서 오는 Spark Connect 세션을 받아들일 준비가 됐어요. 다음 섹션에서는 클라이언트 애플리케이션을 작성할 때 Spark Connect를 사용하는 방법을 살펴볼게요.

대화형 분석에 Spark Connect 사용하기 (Use Spark Connect for interactive analysis)

Spark 세션을 만들 때 Spark Connect를 사용하고 싶다고 지정할 수 있는데, 그 방법에는 다음과 같은 몇 가지가 있어요.

여기서 설명하는 메커니즘 중 하나를 사용하지 않으면, Spark 세션은 Spark Connect를 활용하지 않고 이전처럼 동작해요.

SPARK_REMOTE 환경 변수 설정 (Set SPARK_REMOTE environment variable)

Spark 클라이언트 애플리케이션이 실행되는 클라이언트 머신에서 SPARK_REMOTE 환경 변수를 설정하고 다음 예제처럼 새 Spark 세션을 만들면, 그 세션은 Spark Connect 세션이 돼요. 이 접근 방식으로 Spark Connect 사용을 시작하는 데 코드 변경이 필요 없어요.

터미널 창에서 SPARK_REMOTE 환경 변수를 이전에 컴퓨터에서 시작한 로컬 Spark 서버를 가리키도록 설정하세요:

export SPARK_REMOTE="sc://localhost"

그리고 평소처럼 Spark 셸을 시작하세요:

./bin/pyspark

PySpark 셸은 이제 환영 메시지에 표시된 것처럼 Spark Connect를 사용해 Spark에 연결됐어요:

Client connected to the Spark Connect server at localhost
Spark 세션 생성 시 Spark Connect 지정하기 (Specify Spark Connect when creating Spark session)

Spark 세션을 만들 때 Spark Connect를 사용하고 싶다고 명시적으로 지정할 수도 있어요.

예를 들어 여기서 설명한 것처럼 Spark Connect로 PySpark 셸을 실행할 수 있어요.

PySpark 셸을 Spark Connect로 실행하려면 remote 파라미터를 포함하고 Spark 서버의 위치를 지정하기만 하면 돼요. 이 예제에서는 localhost를 사용해 이전에 시작한 로컬 Spark 서버에 연결해요:

./bin/pyspark --remote "sc://localhost"

그리고 PySpark 셸 환영 메시지가 Spark Connect를 사용해 Spark에 연결했음을 알려주는 것을 볼 수 있어요:

Client connected to the Spark Connect server at localhost

Spark 세션 타입도 확인할 수 있어요. .connect.가 포함되어 있으면 Spark Connect를 사용하고 있는 거예요. 이 예제처럼요:

SparkSession available as 'spark'.
>>> type(spark)
<class 'pyspark.sql.connect.session.SparkSession'>

이제 셸에서 PySpark 코드를 실행해 Spark Connect가 동작하는 것을 확인할 수 있어요:

>>> columns = ["id", "name"]
>>> data = [(1,"Sarah"), (2,"Maria")]
>>> df = spark.createDataFrame(data).toDF(*columns)
>>> df.show()
+---+-----+
| id| name|
+---+-----+
|  1|Sarah|
|  2|Maria|
+---+-----+

Scala 셸의 경우 Ammonite 기반 REPL을 사용해요. 그 외에는 PySpark 셸과 매우 유사해요.

./bin/spark-shell --remote "sc://localhost"

REPL이 성공적으로 초기화되면 인사말 메시지가 나타나요:

Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 4.2.0
      /_/

Type in expressions to have them evaluated.
Spark session available as 'spark'.

기본적으로 REPL은 로컬 Spark 서버에 연결을 시도해요. 셸에서 다음 Scala 코드를 실행해 Spark Connect가 동작하는 것을 확인하세요:

@ spark.range(10).count
res0: Long = 10L
클라이언트-서버 연결 구성 (Configure client-server connection)

기본적으로 REPL은 포트 15002의 로컬 Spark 서버에 연결을 시도해요. 하지만 연결은 이 구성 레퍼런스에 설명된 여러 방식으로 구성될 수 있어요.

SPARK_REMOTE 환경 변수 설정 (Set SPARK_REMOTE environment variable)

SPARK_REMOTE 환경 변수는 클라이언트 머신에서 설정해 REPL 시작 시 초기화되는 클라이언트-서버 연결을 커스터마이즈할 수 있어요.

export SPARK_REMOTE="sc://myhost.com:443/;token=ABCDEFG"
./bin/spark-shell

또는

SPARK_REMOTE="sc://myhost.com:443/;token=ABCDEFG" spark-connect-repl
연결 문자열로 프로그래밍 방식으로 구성 (Configure programmatically with a connection string)

연결은 이 예제처럼 SparkSession#builder를 사용해 프로그래밍 방식으로도 만들 수 있어요:

@ import org.apache.spark.sql.SparkSession
@ val spark = SparkSession.builder.remote("sc://localhost:443/;token=ABCDEFG").getOrCreate()

독립 실행형 애플리케이션에서 Spark Connect 사용하기 (Use Spark Connect in standalone applications)

먼저 pip install pyspark-client==4.2.0로 PySpark를 설치하거나, 패키징된 PySpark 애플리케이션/라이브러리를 빌드하는 경우 setup.py 파일에 다음과 같이 추가하세요:

install_requires=[
'pyspark-client==4.2.0'
]

자신의 코드를 작성할 때는 이 예제처럼 Spark 세션을 만들 때 Spark 서버에 대한 참조를 가진 remote 함수를 포함하세요:

from pyspark.sql import SparkSession
spark = SparkSession.builder.remote("sc://localhost").getOrCreate()

설명을 위해 간단한 Spark Connect 애플리케이션인 SimpleApp.py를 만들어 볼게요:

"""SimpleApp.py"""
from pyspark.sql import SparkSession

logFile = "YOUR_SPARK_HOME/README.md"  # Should be some file on your system
spark = SparkSession.builder.remote("sc://localhost").appName("SimpleApp").getOrCreate()
logData = spark.read.text(logFile).cache()

numAs = logData.filter(logData.value.contains('a')).count()
numBs = logData.filter(logData.value.contains('b')).count()

print("Lines with a: %i, lines with b: %i" % (numAs, numBs))

spark.stop()

이 프로그램은 텍스트 파일에서 'a'를 포함한 줄 수와 'b'를 포함한 줄 수를 세는 것뿐이에요. YOUR_SPARK_HOME을 Spark가 설치된 위치로 바꿔야 한다는 점을 기억하세요.

이 애플리케이션을 일반 Python 인터프리터로 다음과 같이 실행할 수 있어요:

# Use the Python interpreter to run your application
$ python SimpleApp.py
...
Lines with a: 72, lines with b: 39

Scala 애플리케이션/프로젝트의 일부로 Spark Connect를 사용하려면 먼저 올바른 의존성을 포함해야 해요. sbt 빌드 시스템을 예로 들면, build.sbt 파일에 다음 의존성을 추가해요:

libraryDependencies += "org.apache.spark" %% "spark-connect-client-jvm" % "4.2.0"

자신의 코드를 작성할 때는 이 예제처럼 Spark 세션을 만들 때 Spark 서버에 대한 참조를 가진 remote 함수를 포함하세요:

import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder().remote("sc://localhost").getOrCreate()

참고: UDF, filter, map 등과 같은 사용자 정의 코드를 참조하는 연산은 필요한 classfile을 수집·업로드하기 위해 ClassFinder가 등록되어야 해요. 또한 모든 JAR 의존성은 SparkSession#AddArtifact를 사용해 서버에 업로드해야 해요.

예:

import org.apache.spark.sql.connect.client.REPLClassDirMonitor
// Register a ClassFinder to monitor and upload the classfiles from the build output.
val classFinder = new REPLClassDirMonitor(<ABSOLUTE_PATH_TO_BUILD_OUTPUT_DIR>)
spark.registerClassFinder(classFinder)

// Upload JAR dependencies
spark.addArtifact(<ABSOLUTE_PATH_JAR_DEP>)

여기서 ABSOLUTE_PATH_TO_BUILD_OUTPUT_DIR은 빌드 시스템이 classfile을 쓰는 출력 디렉터리이고, ABSOLUTE_PATH_JAR_DEP는 로컬 파일 시스템에서 JAR의 위치예요.

REPLClassDirMonitor는 특정 디렉터리를 모니터링하는 ClassFinder의 제공된 구현이지만, 커스텀 검색·모니터링을 위해 ClassFinder를 확장한 자신만의 클래스를 구현할 수도 있어요.

Spark Connect로 애플리케이션 개발하는 방법과 Spark Connect를 커스텀 기능으로 확장하는 방법에 대한 자세한 내용은 Application Development with Spark Connect를 참고하세요.

클라이언트 애플리케이션 인증 (Client application authentication)

Spark Connect에는 내장 인증이 없지만, 기존 인증 인프라와 매끄럽게 작동하도록 설계됐어요. gRPC HTTP/2 인터페이스는 인증 프록시(authenticating proxies)의 사용을 허용하므로, Spark 자체에 인증 로직을 구현하지 않고도 Spark Connect를 보호할 수 있어요.

지원되는 것 (What is supported)

  • PySpark: Spark 3.4부터 Spark Connect는 DataFrame, Functions, Column을 포함한 대부분의 PySpark API를 지원해요. 하지만 SparkContext와 RDD 같은 일부 API는 지원되지 않아요. 현재 지원되는 API는 API 레퍼런스 문서에서 확인할 수 있어요. 지원되는 API에는 "Supports Spark Connect" 레이블이 붙어 있으므로, 기존 코드를 Spark Connect로 마이그레이션하기 전에 사용 중인 API를 사용할 수 있는지 확인할 수 있어요.
  • Scala: Spark 3.5부터 Spark Connect는 Dataset, functions, Column, Catalog, KeyValueGroupedDataset를 포함한 대부분의 Scala API를 지원해요.
  • 사용자 정의 함수 (User-Defined Functions, UDFs): 지원돼요. 셸에서는 기본적으로, 독립 실행형 애플리케이션에서는 추가 설정 요구 사항과 함께 지원돼요.
  • 스트리밍 API (Streaming API): DataStreamReader, DataStreamWriter, StreamingQuery, StreamingQueryListener를 포함한 스트리밍 API의 대부분이 지원돼요.
  • SparkContext와 RDD 같은 API는 Spark Connect에서 지원되지 않아요.
  • 더 많은 API에 대한 지원이 향후 Spark 릴리스에서 계획돼 있어요.

더 알아보기 (Learn more)