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

    --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}
```
      --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
  ```
      --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 (...)
```

보안, 관측 가능성 및 비용

    --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
 ```
  1. 보안 및 네트워킹 강화

    • GCS 경로, 메타스토어, BigQuery 데이터 세트에 필요한 역할만 부여하는 전용 서비스 계정으로 클러스터를 실행합니다. Private Google Access를 사용하는 제한된 서브넷에 비공개 IP 클러스터를 생성하고 방화벽 규칙을 통해 UI 액세스를 제한합니다.
    • 근거: 최소 권한 및 네트워크 격리는 공격 표면을 줄입니다. 비공개 제어 영역 이그레스는 공개 노출을 방지합니다.
  2. 로깅, 기록 및 알림 계측

    • GCS에 Spark 이벤트 로그를 활성화하고 History Server를 배포합니다. 드라이버/YARN 로그는 보존 정책과 함께 Cloud Logging으로 라우팅합니다. 장기 보류 중인 컨테이너, 반복적인 태스크 실패 또는 과도한 작업 기간에 대한 Monitoring 알림을 추가합니다.
    • 근거: 중앙 집중식 로그는 근본 원인 분석을 지원합니다. 사전 예방적 알림은 스큐, OOM 또는 I/O 성능 저하를 조기에 감지합니다.
  3. 애드혹 및 탄력적 급증에 대비하여 Dataproc Serverless로 선택적 현대화

    • 산발적이거나 탐색적인 Spark SQL 워크로드를 Dataproc Serverless로 이전합니다. 야간 파이프라인은 서버리스에서 완전히 검증될 때까지 임시 클러스터에서 계속 실행합니다.
    • 근거: 서버리스는 클러스터 운영을 제거하고 자동으로 확장되므로 예측 불가능한 부하에 이상적입니다. 기존 워크플로는 최소한의 코드 변경으로 계속 실행됩니다.
  4. 객체 저장소 커미터 및 작은 파일 관리 검증

    • 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.

시험 합격하기 →

Google 찾아보기 →

Related guides

올인원 액세스

하나의 구독. 모든 시험.

모든 플랜은 무제한 답변 검색, 모의고사, AI 해설, 전체 자료 라이브러리를 20개 이상의 언어로 잠금 해제합니다.

월간
24.87
Just €0.83/day
모든 포함:
  • 무제한 답변 검색
  • 무제한 모의고사
  • AI 기반 해설
  • 전체 자료 라이브러리
  • 20개 이상의 언어
  • 주간 콘텐츠 업데이트
  • 보상 및 추천
  • 우선 지원
무료 체험 시작

신용카드 필요 없음*

최고의 가치
12개월
179.87
Just €0.49/daySave 40%
모든 포함:
  • 무제한 답변 검색
  • 무제한 모의고사
  • AI 기반 해설
  • 전체 자료 라이브러리
  • 20개 이상의 언어
  • 주간 콘텐츠 업데이트
  • 보상 및 추천
  • 우선 지원
무료 체험 시작

신용카드 필요 없음*

✓ 무료 플랜 포함 · ✓ 언제든지 취소 가능 · ✓ 모든 플랜은 전체 제품을 잠금 해제합니다