Spark DataFrame API란 무엇인가요?

Apache Spark는 정형, 반정형, 비정형 데이터를 대규모로 처리하도록 설계된 분산 처리 엔진입니다. 개발자는 Spark API를 사용하여 엔진과 상호작용하고 분산 데이터 작업을 정의합니다.


최신 Spark 개발의 핵심은 Spark DataFrame API입니다. DataFrame은 관계형 데이터베이스의 테이블이나 스프레드시트와 마찬가지로 분산 데이터를 이름이 지정된 열로 구성하고 처리합니다. 구조화된 선언적 인터페이스를 제공하여 대규모 데이터 세트를 분석하는 동시에 백그라운드에서 물리적 실행을 자동으로 최적화합니다.

Spark API 아키텍처 및 실행

Spark API가 분산 네트워크에서 코드를 실행하는 방식을 이해하려면 런타임 아키텍처의 기본 구성요소를 이해하는 것이 중요합니다.


  • 드라이버: Spark 애플리케이션의 중앙 코디네이터입니다. 코드를 읽고, 선언적 작업을 논리적 실행 계획으로 변환하고, 워커 노드 전반에 걸쳐 작업을 예약합니다.
  • 클러스터 관리자: 클러스터 전반에 CPU, 메모리, 네트워킹 리소스를 프로비저닝하는 리소스 할당자(예: Kubernetes 또는 독립형 리소스 관리자)입니다.
  • 실행자: 클러스터 노드에서 실행되는 작업자 인스턴스입니다. 드라이버로부터 태스크 명령을 수신하고, 메모리에서 로컬로 데이터 처리 작업을 실행하고, 결과를 반환하거나 스토리지에 기록합니다.
  • DAG, 단계, 태스크: Spark는 데이터 변환을 한 줄씩 실행하지 않습니다. 대신 드라이버는 DataFrame 코드를 물리적 연산자의 방향성 비순환 그래프(DAG)로 컴파일합니다. 드라이버는 이러한 연산자를 데이터 셔플링 경계로 구분되는 단계라는 광범위한 실행 단계로 그룹화하고, 이러한 단계를 병렬 실행을 위해 실행자에게 분산되는 개별 태스크로 세분화합니다.

Spark RDD와 DataFrame의 차이점은 무엇인가요?

분산 데이터 파이프라인을 설계할 때 개발자는 두 가지 기본 프로그래밍 인터페이스 중에서 선택합니다.


  • Spark RDD(탄력적 분산 데이터 세트): 원래의 하위 수준 Spark API입니다. RDD는 클러스터 전반에 분산된 JVM 객체의 변경 불가능하고 내결함성이 있는 컬렉션을 나타냅니다. RDD에는 임의의 Java 객체가 포함되어 있으므로 실행 엔진이 내부 구조를 검사할 수 없습니다. 이로 인해 개발자는 하위 수준 변환 로직을 수동으로 작성하고 조정해야 하므로 상당한 가비지 컬렉션 오버헤드와 느린 직렬화가 발생하는 경우가 많습니다.
  • Spark DataFrame: 구조화된 데이터 처리를 위한 최신 표준입니다. DataFrame은 엄격한 스키마로 정의된 행과 열로 데이터를 구성합니다. 실행 엔진은 데이터 세트의 데이터 유형과 구조를 이해하므로 코드를 실행하기 전에 자동으로 최적화할 수 있습니다.

RDD 대신 DataFrame을 선택해야 하는 이유

거의 모든 데이터 엔지니어링 및 데이터 과학 사용 사례에서 DataFrame API가 선호됩니다. RDD는 커스텀 Java 객체를 사용하여 원시 비정형 데이터(예: 바이너리 파일 또는 미디어 스트림)를 조작해야 하는 시나리오를 위해 예약되어 있습니다. DataFrame은 더 적은 코드를 필요로 하고, 자동으로 더 빠르게 실행되며, 오프힙 바이너리 메모리 관리를 활용하여 JVM 가비지 컬렉션 병목 현상을 우회합니다.

데이터팀에서 Spark API를 사용하는 방식

Spark DataFrame API는 대량의 데이터를 처리, 정리, 준비하는 계산 집약적인 작업을 간소화합니다.


  • 데이터 엔지니어: 엔지니어는 DataFrame API를 사용하여 복원력과 내결함성이 뛰어난 데이터 파이프라인을 빌드합니다. 원시 데이터를 추출하고, 스키마를 적용하고, 구조화된 변환을 실행하고, 일치하는 테이블을 데이터 레이크 또는 클라우드 데이터 웨어하우스에 작성합니다. 선언적 문법을 사용하면 코드 복잡성과 유지보수 오버헤드를 최소화하면서 매일 테라바이트 규모의 데이터를 처리할 수 있습니다.
  • 데이터 과학자: 데이터 과학자는 DataFrame API를 사용하여 단일 머신의 메모리 한도를 초과하는 대규모 데이터 세트를 탐색하고 준비합니다. 데이터 파티션을 워커 노드에 분산함으로써 탐색적 데이터 분석을 실행하고, null 값을 정리하고, 특성을 대규모로 엔지니어링하여 인사이트 도출 시간을 단축할 수 있습니다.

DataFrame API를 기반으로 빌드된 핵심 라이브러리

Spark DataFrame API는 Spark의 고급 처리 라이브러리를 위한 프로그래매틱 기반 역할을 합니다.


  • Spark SQL: 이 모듈을 사용하면 Spark 데이터 세트에서 직접 ANSI SQL 쿼리를 실행할 수 있습니다. DataFrame을 임시 뷰로 쿼리하여 선언적 Python 또는 Scala 코드를 표준 SQL 쿼리와 원활하게 혼합할 수 있습니다.
  • Structured Streaming: 이 엔진은 실시간 데이터의 연속 스트림을 처리합니다. 정적 일괄 처리와 동일한 DataFrame API 명령어를 사용하며 마이크로 일괄 처리, 내결함성, 상태 관리를 자동으로 처리합니다.
  • MLlib: Spark의 분산 머신러닝 라이브러리는 DataFrame을 사용하여 학습 데이터 세트를 관리하고 준비하며, 분류, 회귀, 클러스터링, 협업 필터링을 위한 확장 가능한 기본 제공 알고리즘을 제공합니다.

Spark DataFrame API의 이점

선언적 최적화

DataFrame은 Catalyst Optimizer라는 기본 제공 쿼리 옵티마이저를 사용합니다. DataFrame 코드를 작성하면 옵티마이저가 필터 푸시다운, 프로젝션 가지치기, 조인 재정렬을 적용하는 물리적 실행 계획을 자동으로 다시 작성하여 수동 조정 없이 최대한 빠르게 실행합니다.

언어 유연성

팀에서 Python(PySpark), Scala, Java, R 중 어떤 언어로 코드를 작성하든 DataFrame API는 동일한 성능을 보장합니다. 기본 실행 계획은 동일한 JVM 독립 실행 레이어 내에서 컴파일되고 최적화됩니다.

효율적인 메모리 관리

DataFrame은 Tungsten 실행 엔진을 활용하여 데이터를 고도로 압축된 오프힙 바이너리 형식으로 저장합니다. 이렇게 하면 JVM 객체 생성 오버헤드가 제거되고 가비지 컬렉션 일시중지로 인해 실행자 스레드가 정지되는 것을 방지할 수 있습니다.

Google Cloud에서 Spark DataFrame API를 사용하는 방법

기존에는 오픈소스 Spark를 실행하려면 수동 클러스터 구성, 소프트웨어 버전 관리, 복잡한 네트워크 피어링이 필요했습니다. Google Cloud는 Managed Service for Apache Spark를 제공하여 분산 실행을 완전 관리형 엔터프라이즈급 데이터 플랫폼으로 전환함으로써 이 과정을 간소화합니다.


데이터팀은 다음 기능을 사용하여 Google Cloud에서 Spark DataFrame 워크로드를 실행합니다.


  • 서버리스 및 관리형 클러스터 배포: Google Cloud를 사용하면 운영 요구사항에 맞는 실행 모델을 선택할 수 있습니다. 서버리스 모드에서 DataFrame 코드를 실행하여 일괄 작업을 직접 제출할 수 있습니다. 리소스가 자동으로 확장되므로 정확한 런타임 시간(초)에 대해서만 비용을 지불하면 됩니다. 또는 지속적인 워크로드를 위해 고도로 맞춤설정 가능한 영구 관리형 클러스터를 배포할 수도 있습니다.
  • 최적화된 BigQuery 연결: Spark-BigQuery 커넥터는 기본 Apache Arrow 형식으로 BigQuery 데이터를 직접 사용함으로써 JVM 행 지향 전환을 우회합니다. 데이터 과학자는 통합 노트북 인터페이스를 사용하여 환경을 전환하지 않고도 동일한 거버넌스 데이터 세트에서 PySpark DataFrame 코드와 SQL 쿼리를 실행할 수 있습니다.
  • 벡터화된 네이티브 C++ 실행: Google Cloud에서 Spark 워크로드를 실행할 때 서버리스 일괄 또는 관리형 클러스터에 Lightning Engine을 사용 설정할 수 있습니다. 이 네이티브 C++ 쿼리 실행 엔진은 Velox와 Gluten을 사용하여 물리적 쿼리 계획을 네이티브 명령어로 직접 컴파일합니다. JVM Volcano 반복자 모델과 가비지 컬렉션 병목 현상을 우회하여 코드 변경 없이 DataFrame 및 Spark SQL 워크로드를 최대 4.9배 가속화합니다.


다음 단계 수행

$300의 무료 크레딧과 20여 개의 항상 무료 제품으로 Google Cloud에서 빌드하세요.

Google Cloud