Google PDE: Spark, Dataproc 및 분산 데이터 처리 — 학습 가이드
다음의 일부입니다: Google Professional Data Engineer — 학습 가이드. 검증된 답안으로 연습하기: Google 시험 허브, 또는 다음에서 시간 제한 모의고사 풀기: ExamRoll.io.
개요
Google Cloud Dataproc의 Apache Spark는 분산 데이터 처리를 위한 관리형 탄력적 플랫폼을 제공합니다. 제어 요구사항, 런타임 가변성, 관리 오버헤드에 따라 상시 실행(long-running) 또는 임시(ephemeral) Dataproc 클러스터와 Dataproc Serverless for Spark 중에서 선택할 수 있습니다. Spark는 복원력 있는 추상화(RDD), 관계형 API(DataFrame 및 Spark SQL), 대규모 반복 및 배치 ETL에 최적화된 내결함성 DAG 실행 엔진을 제공합니다. Google Cloud에서는 Cloud Storage가 HDFS를 대체하여 내구성 있고 저렴한 스토리지를 제공하고, BigQuery 커넥터는 직접적인 분석 오프로드를 지원하며, Dataproc Metastore는 스키마 관리를 중앙화합니다. 효과적인 솔루션은 스토리지와 컴퓨팅 수명 주기를 일치시키고, 워크로드에 맞게 Spark를 조정하며, 관측 가능성을 계측하고, 최소 권한 및 네트워크 격리를 통해 보안을 적용합니다.
Dataproc 아키텍처: 클러스터, 서버리스, 스토리지, Metastore
- 클러스터 유형 및 노드 역할
- 기본(마스터) 노드는 YARN, HDFS NameNode(사용 시), Spark 드라이버 UI를 호스팅합니다. HA 모드는 여러 기본 노드를 사용합니다.
- 작업자 노드는 실행기(executor)와 HDFS DataNode(사용 시)를 실행합니다.
- 보조/추가 작업자는 일반적으로 선점형/스팟(spot) VM으로, HDFS 역할 없이 탄력적이고 저렴한 용량을 제공합니다.
- 이미지는 OS와 구성요소 버전을 번들로 제공합니다(예: 2.1-debian11, 2.2-ubuntu20). 이미지 버전을 고정하여 Spark/Hadoop 호환성을 제어하고 신중하게 업그레이드해야 합니다.
- Component Gateway는 UI(Spark History Server, YARN RM)를 HTTPS를 통해 안전하게 게시합니다.
- Dataproc Serverless for Spark
- 클러스터 프로비저닝이 필요 없고, 자동 확장(autoscaling)이 적용되며, 실행기와 드라이버에 대해 초 단위로 과금됩니다. 산발적이거나 폭증하는(bursty) 작업 또는 운영 오버헤드를 최소화해야 할 때 이상적입니다.
- 단점: 클러스터보다 하위 수준의 조정 옵션이 적고, 작업 시작 지연 시간이 준비된(warm) 클러스터보다 길 수 있습니다. 문제 해결을 위해 서버리스 측정항목과 이벤트 로그를 사용합니다.
- 자동 확장(Autoscaling)
- 클러스터 자동 확장 정책은 YARN/Spark 측정항목과 유예 기간(cooldown)을 기반으로 작업자를 추가/제거하며, 기본 작업자 그룹과 보조 작업자 그룹을 별도로 조정합니다.
- 서버리스 자동 확장은 서비스에서 관리합니다. 최상의 확장을 위해 파티션 병렬 처리가 가능하도록 설계하고 직렬화된 병목 현상을 피해야 합니다.
- 스토리지 및 커넥터
- 기본 스토리지 시스템(system-of-record)으로 Google Cloud Storage(GCS)를 사용하는 것이 좋습니다. GCS는 컴퓨팅과 스토리지를 분리하고, 영구 디스크 비용을 절감하며, 클러스터 수명 주기와 무관하게 데이터를 보존합니다.
- GCS 커넥터(gs://)는 Hadoop/Spark와 통합됩니다. 객체 스토어에 쓸 때는 커밋 프로토콜을 사용합니다. GCS에서 이름 변경 오버헤드를 줄이고 작업 커밋 속도를 높이려면 FileOutputCommitter 알고리즘 v2를 설정하십시오:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- 열 제거(column pruning) 및 조건자 푸시다운(predicate pushdown)과 함께 Parquet/ORC를 사용하십시오. 작은 파일은 압축(compaction)을 통해 관리하여 효율적인 스캔을 위해 파일당 128–512MiB를 목표로 합니다.
- Hive metastore
- 스키마와 테이블 메타데이터를 Dataproc Metastore(관리형 Apache Hive Metastore) 또는 Cloud SQL 기반 metastore에 중앙 집중화하여 여러 클러스터에서 카탈로그를 공유합니다.
- 내구성을 위해 GCS를 가리키는 외부 테이블을 사용하고, 날짜/시간별로 파티션을 나누어 스캔 비용을 제한합니다.
- 작업, 초기화, 워크플로
- spark, pyspark, spark-sql 또는 hadoop 작업을 제출합니다. 초기화 작업을 사용하면 클러스터 생성 시 추가 라이브러리나 에이전트(예: 커넥터, Python 라이브러리)를 설치할 수 있습니다.
- 워크플로 템플릿은 다단계 파이프라인을 매개변수화합니다. 워크플로별로 임시 클러스터를 생성한 후 해체할 수 있습니다. 이를 통해 격리 수준을 높이고 유휴 비용을 절감할 수 있습니다.
- 배치 ETL에는 임시 클러스터를 권장합니다. 데이터와 metastore는 클러스터 외부(GCS, Dataproc Metastore, BigQuery)에 저장됩니다.
- BigQuery 통합
- Spark BigQuery 커넥터는 BigQuery를 직접 읽고 씁니다. 처리량을 높이려면 BigQuery Storage Read API를, 지연 시간이 짧고 정확히 한 번(exactly-once) 스트리밍 삽입을 위해서는 Write API를 고려하십시오.
- 테이블 유지 관리를 위해 BigQuery에서 다운스트림 MERGE/파티션 덮어쓰기를 수행하여 로드를 원자적으로 완료합니다.
### Spark 모델, 성능 튜닝 및 안정성
- API 및 실행
- RDD: 저수준, 불변(immutable), Scala/Java에서 타입-세이프(type-safe)함. 파티셔닝과 영속성(persistence)을 직접 제어.
- DataFrame/Dataset: 관계형, Catalyst에 의해 최적화됨. 쿼리 최적화 및 코드 생성 기능 때문에 ETL 작업에 사용하는 것을 권장.
- Transformation(map, filter, join 등)은 지연(lazy) 평가됨. Action(count, collect, save 등)이 실행을 트리거함. Spark는 셔플(shuffle)에 의해 분리된 스테이지(stage)들의 DAG를 구축하며, 태스크(task)는 파티션별로 실행됨.
- 파티셔닝 및 셔플
- 입력 파티셔닝: 모든 코어를 활용할 수 있도록 충분한 파티션을 확보. 전체 executor 코어 수의 2–4배로 시작. RDD의 경우 spark.default.parallelism으로, DataFrame의 경우 리더(reader) 옵션으로 제어.
- 셔플 파티션: 기본값 200은 종종 파티션 수가 너무 적거나 많게 설정됨. 다음과 같이 튜닝:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- 와이드 트랜스포메이션(wide transform) 후 파티션당 약 100–256 MiB를 목표로 함. 너무 작으면 스케줄러 오버헤드가 발생하고, 너무 크면 executor OOM 위험이 있음.
- 셔플은 join, groupBy, orderBy 작업에서 가장 큰 비용을 차지함. 충분한 executor 메모리와 디스크를 확보하고, 클러스터에서 셔플이 많은 경우 로컬 SSD 사용을 고려.
- 스큐(Skew) 및 조인 전략
- 스큐 탐지(긴 태스크 실행 시간, 큰 파티션 크기). 완화 방법:
- 작은 테이블을 브로드캐스트하여 셔플 방지:
- 스큐 탐지(긴 태스크 실행 시간, 큰 파티션 크기). 완화 방법:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- 핫 파티션(hot partition)에 키 솔팅(salt key) 적용, 맵 사이드 사전 집계(map-side pre-aggregation) 적용, 조기 필터링.
- Adaptive Query Execution(AQE)을 활성화하여 셔플 후 파티션을 병합하고 스큐 조인(skewed join) 처리:
--conf spark.sql.adaptive.enabled=true
```
- 캐싱, 체크포인팅 및 리니지
- 재사용되는 핫 중간(hot intermediate) DataFrame은 신중하게 캐시함. OOM을 피하기 위해 MEMORY_AND_DISK를 선호.
- 긴 리니지(lineage)는 GCS나 HDFS에 체크포인트하여 장애 발생 시 재계산 범위를 제한.
- Executor 및 동적 할당
- 병렬성과 GC 오버헤드의 균형을 맞추기 위해 executor 크기를 적절히 조절:
- Executor당 코어 수: I/O/CPU 작업의 균형을 위해 2–5개. 코어 수가 적을수록 GC 중단 시간이 줄어듦.
- 메모리 오버헤드: 와이드 셔플을 위해 spark.yarn.executor.memoryOverhead 설정.
- 클러스터에서 외부 셔플 서비스와 함께 동적 할당을 활성화하여 워크로드에 따라 executor를 확장:
- 병렬성과 GC 오버헤드의 균형을 맞추기 위해 executor 크기를 적절히 조절:
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- 배치 ETL을 위한 내결함성 패턴
- 멱등성(Idempotent) 쓰기: 임시/스테이징 경로에 쓰고, 디렉터리 수준 커밋으로 원자적으로 승격. BigQuery의 경우 스테이징 테이블에 쓰고 MERGE 사용:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- 증분 처리: ingestion_date 파티션에 워터마크 기반 필터링 사용. 재처리를 방지하기 위해 GCS에 처리된 매니페스트(processed-manifest)를 유지.
- 데드-레터(Dead-letter) 처리: 파싱/유효성 검사 오류 시, 잘못된 레코드를 진단 정보와 함께 격리 경로/테이블로 분기. 엄격한 스키마 적용 및 내장 DLQ가 필요하면 Dataflow를 고려. Spark에서는 레코드별 try/catch와 별도의 싱크(sink)를 구현.
보안, 관측 가능성 및 비용
- ID 및 액세스
- 최소 권한 IAM을 사용하는 전용 서비스 계정으로 클러스터와 작업을 실행합니다. 필요한 역할만 부여합니다. 예:
- 인스턴스 서비스 계정에 roles/dataproc.worker
- GCS I/O 경로에 roles/storage.objectViewer 또는 objectAdmin
- 대상 데이터 세트에 roles/bigquery.dataEditor
- Dataproc Serverless의 경우, 작업별 서비스 계정을 사용하여 액세스 범위를 지정합니다.
- 최소 권한 IAM을 사용하는 전용 서비스 계정으로 클러스터와 작업을 실행합니다. 필요한 역할만 부여합니다. 예:
- 네트워크 격리 및 암호화
- VPC 서브넷에서 비공개 IP 클러스터를 사용하고, 방화벽으로 마스터 UI를 제한하며, 공개 이그레스 없이 GCS/BigQuery에 액세스할 수 있도록 Private Google Access를 활성화합니다.
- 중앙 집중식 제어를 위해 Shared VPC 프로젝트에 클러스터를 배치합니다. 선택적으로 Dataproc에서 Kerberos를 활성화하여 클러스터 내 인증을 수행할 수 있습니다.
- CMEK로 저장 데이터 암호화: GCS 버킷, Persistent Disks, Dataproc Metastore, BigQuery에 CMEK를 구성합니다. 전송 중 데이터는 기본적으로 TLS를 사용합니다.
- 로깅, 기록 및 측정항목
- GCS에 Spark 이벤트 로그를 활성화하고 History Server를 배포합니다:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc은 드라이버 및 YARN 로그를 Cloud Logging으로 스트리밍합니다. 보관/포렌식을 위해 싱크로 내보냅니다.
- Cloud Monitoring 측정항목으로 모니터링: YARN 보류 중인 컨테이너, CPU, 메모리, HDFS 상태(사용 시), GCS 처리량. 장기적인 스테이지 재시도, 실행기 손실, 추측 실행 급증 시 알림을 설정합니다.
- 장애 분석: 일반적인 원인으로는 스큐로 인한 스트래글러(straggler), 셔플 중 실행기 OOM, 객체 저장소 커밋 실패, 선점형/스팟 노드 손실 등이 있습니다. 재시도 횟수는 신중하게 늘려야 합니다. 과도한 재시도는 비용과 지연을 증폭시킬 수 있습니다.
- 비용 최적화
- 임시 클러스터 또는 Dataproc Serverless를 사용하여 유휴 비용을 방지하고, 영구 디스크 사용을 최소화하기 위해 데이터는 GCS에 보관합니다.
- 선점형/스팟 보조 작업자를 추가하여 피크 수요를 흡수합니다. 손실된 노드의 태스크는 재시도되므로 재계산을 염두에 두고 설계합니다. 마스터 노드는 선점형 노드에 배치하지 마십시오.
- 적절한 머신 유형을 선택하고 자동 확장을 사용하여 큐가 비어 있을 때 용량을 축소합니다. 스캔 비용과 CPU를 줄이기 위해 파티션 프루닝(partition pruning)이 가능한 Parquet/ORC를 선호합니다.
- 출력을 압축하여 작은 파일 생성을 피합니다. 파일 수가 적고 크기가 클수록 메타데이터 오버헤드와 작업 실행 시간이 줄어듭니다.
- 짧고 주기적인 작업(예: 매주 30분짜리 Spark ETL)의 경우, 선점형 작업자나 서버리스가 종종 최상의 비용 프로필을 제공합니다.
#### 실용적인 문제 시나리오
Acme Retail은 다운스트림 분석에 데이터를 제공하는 야간 Spark 및 Hive ETL을 실행하는 30개 노드의 온프레미스 Hadoop 클러스터를 마이그레이션하려고 합니다. 이들은 기존 작업을 최소한의 변경으로 재사용하고, 클러스터를 상시 관리하는 것을 피하며, 클러스터 수명 주기 이후에도 데이터를 영구 보존하고, 스토리지 비용을 절감하고자 합니다.
접근 방식:
1) 관리형 서비스에 데이터 및 메타데이터 저장
- 모든 원시 데이터와 큐레이션된 데이터를 파티셔닝(예: dt=YYYY-MM-DD)된 Parquet 형식을 사용하여 Cloud Storage에 저장합니다.
- 근거: GCS는 내구성이 뛰어나고 저렴하며, 컴퓨팅과 스토리지를 분리하므로 임시 클러스터와 서버리스 작업이 영구 디스크 없이 실행될 수 있습니다. 파티셔닝된 Parquet는 조건자 푸시다운(predicate pushdown)과 효율적인 스캔을 가능하게 합니다.
2) Dataproc Metastore로 카탈로그 중앙화
- Hive 메타스토어를 Dataproc Metastore로 마이그레이션합니다. GCS 경로를 참조하는 외부 Hive 테이블을 생성하고 기존 스키마/파티션 로직을 유지합니다.
- 근거: 관리형 메타스토어를 사용하면 여러 임시 클러스터와 서버리스 작업이 HA MySQL/PostgreSQL 인스턴스를 실행하지 않고도 테이블 정의를 공유할 수 있습니다.
3) 배치 ETL에는 임시 Dataproc 클러스터를, 오케스트레이션에는 워크플로 템플릿을 사용
- 필요한 이미지(예: 2.1-debian11)로 클러스터를 생성하고, Spark 작업(spark-sql 및 pyspark)을 실행한 후, 완료 시 클러스터를 삭제하는 워크플로 템플릿을 정의합니다. 초기화 작업을 추가하여 사용자 정의 라이브러리를 설치합니다.
- 근거: 임시 클러스터는 유휴 비용을 제거하고 작업 종속성을 격리합니다. 워크플로 템플릿은 반복성과 매개변수화(날짜, 입력 경로)를 제공합니다.
4) 자동 확장 및 선점형 작업자 활성화
- 소규모 핵심 작업자 그룹과 대규모 선점형 보조 작업자 풀로 구성된 자동 확장 정책을 연결합니다. 실행 후 신속하게 축소되도록 재사용 대기시간(cooldown)을 조정합니다.
- 근거: 핵심 작업자는 클러스터 안정성을 유지하고, 선점형 작업자는 셔플과 광범위한 변환을 저렴한 비용으로 흡수합니다. Spark/YARN 재시도는 선점 시 손실된 태스크를 처리합니다.
5) Spark BigQuery 커넥터를 통해 BigQuery와 통합
- 차원/팩트 로드의 경우, Spark 결과를 스테이징 BigQuery 테이블에 쓰고, MERGE 문을 실행하여 대상을 원자적으로 업데이트합니다. 직접 덮어쓰기가 안전한 경우, 파티션 덮어쓰기 모드를 사용하여 파티셔닝된 테이블에 씁니다.
- 근거: BigQuery는 대규모 분석 및 BI를 제공합니다. 스테이징+MERGE 방식은 배치 Spark에서 트랜잭션과 유사한 업서트(upsert)를 생성하여 다운스트림 불일치를 줄입니다.
6) 성능 및 안정성을 위한 Spark 튜닝
- 실행기 코어 수에 비례하여 셔플 파티션을 설정하고 AQE를 활성화합니다:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- 작은 차원 테이블에는 브로드캐스트 조인을 사용하고, 안정성을 위해 긴 리니지(lineage)는 GCS에 체크포인팅합니다.
- 근거: 적절한 파티셔닝은 스큐와 스케줄러 오버헤드를 줄입니다. AQE는 런타임에 데이터 프로필에 적응합니다. 체크포인팅은 장애 발생 후 재계산 범위를 제한합니다.
보안 및 네트워킹 강화
- GCS 경로, 메타스토어, BigQuery 데이터 세트에 필요한 역할만 부여하는 전용 서비스 계정으로 클러스터를 실행합니다. Private Google Access를 사용하는 제한된 서브넷에 비공개 IP 클러스터를 생성하고 방화벽 규칙을 통해 UI 액세스를 제한합니다.
- 근거: 최소 권한 및 네트워크 격리는 공격 표면을 줄입니다. 비공개 제어 영역 이그레스는 공개 노출을 방지합니다.
로깅, 기록 및 알림 계측
- GCS에 Spark 이벤트 로그를 활성화하고 History Server를 배포합니다. 드라이버/YARN 로그는 보존 정책과 함께 Cloud Logging으로 라우팅합니다. 장기 보류 중인 컨테이너, 반복적인 태스크 실패 또는 과도한 작업 기간에 대한 Monitoring 알림을 추가합니다.
- 근거: 중앙 집중식 로그는 근본 원인 분석을 지원합니다. 사전 예방적 알림은 스큐, OOM 또는 I/O 성능 저하를 조기에 감지합니다.
애드혹 및 탄력적 급증에 대비하여 Dataproc Serverless로 선택적 현대화
- 산발적이거나 탐색적인 Spark SQL 워크로드를 Dataproc Serverless로 이전합니다. 야간 파이프라인은 서버리스에서 완전히 검증될 때까지 임시 클러스터에서 계속 실행합니다.
- 근거: 서버리스는 클러스터 운영을 제거하고 자동으로 확장되므로 예측 불가능한 부하에 이상적입니다. 기존 워크플로는 최소한의 코드 변경으로 계속 실행됩니다.
객체 저장소 커미터 및 작은 파일 관리 검증
- FileOutputCommitter 알고리즘을 v2로 설정하고, 쓰기 전에 repartition/coalesce를 통해 출력을 파일당 256–512 MiB로 압축합니다.
- 근거: 객체 저장소에는 원자적 이름 바꾸기 기능이 없습니다. 최적화된 커미터는 복사/이름 바꾸기 오버헤드를 줄입니다. 압축은 성능 및 비용 측면에서 작은 파일 문제를 완화합니다.
이 설계는 최소한의 리팩토링으로 기존 Spark 및 Hive 작업을 재사용하고, GCS에서 데이터 내구성을 보장하며, 스키마를 중앙 집중화하고, 보안 영향 반경을 억제하며, 강력한 관측 가능성을 제공하고, 임시 클러스터, 자동 확장, 선점형 용량 및 서버리스 실행의 선별적 사용을 통해 비용을 최적화합니다.
이 문제 연습하기 → · ExamRoll.io에서 시간 제한 연습 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
시험 합격하기 →