Spark 선언적 파이프라인
Spark 선언적 파이프라인(SDP)은 AWS Glue 6.0에서 배치 및 스트리밍 데이터 파이프라인을 빌드하기 위한 선언적 프레임워크입니다. SDP에서는 SQL 또는 Python을 사용하여 데이터의 모양을 정의할 수 있으며, 프레임워크는 실행 계획을 자동으로 결정하고, 데이터세트 간의 종속성을 해결하고, 독립 브랜치를 병렬로 실행합니다.
SDP는 읽기, 쓰기, 카탈로그 등록, 실행 순서를 위한 보일러플레이트 코드를 제거하여 파이프라인 개발을 간소화합니다. 프레임워크가 파이프라인 인프라를 처리하는 동안 사용자는 비즈니스 혁신에 집중할 수 있습니다.
SDP는 AWS Glue 버전 6.0 이상에서 사용할 수 있습니다.
SDP 개념
파이프라인은 YAML 매니페스트 파일(spark-pipeline.yml)과 하나 이상의 SQL 또는 Python 변환 파일로 구성됩니다. SDP는 자동으로 다음을 수행합니다.
-
테이블 참조에서 DAG를 추론하여 종속성을 해결합니다.
-
수동 오케스트레이션 없이 실행 순서를 결정합니다.
-
최대 처리량을 위해 독립 브랜치를 병렬로 실행합니다.
-
체크포인트를 통해 스트리밍 테이블의 증분 상태를 관리합니다.
-
구체화 시 카탈로그에 출력 테이블을 등록합니다.
데이터세트 유형
SDP에서는 세 가지 데이터세트 유형을 사용할 수 있습니다.
- 스트리밍 테이블
-
마지막 실행 이후 새 데이터만 처리합니다. 체크포인트를 사용하여 작업 실행 간에 상태를 유지합니다. 수집, 이벤트 스트림, IoT 데이터, 변경 데이터 캡처, 추가 전용 소스에 스트리밍 테이블을 사용합니다.
- 구체화된 뷰
-
실행할 때마다 데이터세트를 완전히 다시 계산합니다. 출력은 항상 소스 데이터의 현재 상태를 반영합니다. 집계, 조인, 요약 분석, 보고서에 구체화된 뷰를 사용합니다.
- 임시 뷰
-
세션 범위로 제한되며 지속되거나 카탈로그화되지 않습니다. 중간 변환 및 스테이징 로직에 임시 뷰를 사용합니다.
중요
구체화된 뷰는 항상 현재 버전에서 전체 재계산을 수행합니다. 증분 새로 고침은 지원하지 않습니다. 증분 워크로드에는 스트리밍 테이블을 사용하세요.
사전 조건
SDP를 사용하려면 다음이 필요합니다.
-
AWS Glue 버전 6.0
-
파이프라인 스토리지(체크포인트, 메타데이터)를 위한 Amazon S3 위치
-
Data Catalog 통합의 경우(선택 사항):
--enable-glue-datacatalog를true로 설정합니다. 또는 Spark 구성을 통해 직접 카탈로그 설정을 구성할 수 있습니다. -
영구 테이블 스토리지의 경우:
spark.sql.warehouse.dir을 Amazon S3 경로로 설정하거나 파이프라인 YAML에서database:필드를 설정하고 AWS Glue 데이터베이스에 Amazon S3 경로로 구성된LocationUri가 있는지 확인합니다. -
스트리밍 테이블을 사용한 교차 실행 증분 처리의 경우 스트리밍 테이블의 데이터 및 체크포인트 상태가 Amazon S3에 지속되어야 합니다. 예를 들어 Hive 또는 AWS Glue 관리형(비 Iceberg) 테이블은 데이터베이스
LocationUri를 Amazon S3 경로로 설정해야 하는 반면 Apache Iceberg 테이블은 테이블 메타데이터 자체를 관리합니다.
중요
Hive 또는 AWS Glue 관리형(비 Iceberg) 테이블에 대한 파이프라인 YAML에서 database: 필드를 사용하는 경우 해당 AWS Glue 데이터베이스의 LocationUri가 Amazon S3 경로로 설정해야 합니다. LocationUri는 관리형 스트리밍 테이블(및 해당 _spark_metadata 로그)을 Amazon S3에 배치하며, 이를 통해 해당 테이블에 대한 지속적 교차 실행 증분 처리를 할 수 있게 됩니다. catalog-impl=GlueCatalog와 함께 데이터 카탈로그를 사용하는 Iceberg 테이블(옵션 1)에는 데이터베이스 LocationUri가 필요하지 않습니다. 명시적 Amazon S3 위치로 데이터베이스를 생성하거나 업데이트합니다.
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
파이프라인 생성
SDP 파이프라인을 생성하려면 다음 단계를 완료합니다.
1단계: 파이프라인 YAML 생성
spark-pipeline.yml라는 이름의 파일을 만듭니다.
name: my_analytics_pipeline catalog: spark_catalog database: analytics_db storage: s3://my-bucket/pipeline-storage/ libraries: - glob: include: transformations/** configuration: spark.sql.shuffle.partitions: "4"
다음 표에서는 파이프라인 YAML 필드를 설명합니다.
| Field | 필수 | 설명 |
|---|---|---|
name |
예 | 파이프라인의 이름입니다. |
catalog |
아니요 | 사용할 카탈로그입니다. 기본값은 spark_catalog입니다. |
database |
아니요 | 출력 테이블의 AWS Glue 데이터베이스입니다. LocationUri 요구 사항은 사전 조건을 참조하십시오. |
storage |
예 | 파이프라인 체크포인트 및 메타데이터에 대한 Amazon S3 경로입니다. |
libraries |
예 | 포함할 변환 파일의 Glob 패턴입니다. |
configuration |
아니요 | Spark 구성 속성입니다. |
2단계: 변환 쓰기
transformations/ 디렉터리에 변환 파일을 생성합니다. 동일한 파이프라인에서 SQL, Python 또는 둘 다를 사용할 수 있습니다.
SQL 예제(transformations/silver.sql):
CREATE MATERIALIZED VIEW silver_sales AS SELECT *, UPPER(region) as clean_region FROM bronze_sales WHERE amount > 0; CREATE MATERIALIZED VIEW gold_summary AS SELECT clean_region, COUNT(*) as order_count, SUM(amount) as total_revenue FROM silver_sales GROUP BY clean_region;
Python 예제(transformations/bronze.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() @dp.materialized_view(comment="Raw sales data from S3") def bronze_sales() -> DataFrame: return spark.read.format("csv").option("header", "true") \ .option("inferSchema", "true") \ .load("s3://source-bucket/raw-data/sales/")
Python 스트리밍 테이블 예제(transformations/events.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() dp.create_streaming_table( "streaming_events", comment="Incremental event ingestion", schema="event_id STRING, event_type STRING, timestamp LONG, payload STRING" ) @dp.append_flow(target="streaming_events") def ingest_events() -> DataFrame: return ( spark.readStream.format("json") .schema("event_id STRING, event_type STRING, timestamp LONG, payload STRING") .load("s3://source-bucket/events/") )
참고
Python에서 스트리밍 테이블을 사용하려면 dp.create_streaming_table()을 @dp.append_flow(target=...)와 함께 사용합니다.
3단계: Amazon S3에 업로드
파이프라인 파일을 Amazon S3에 다음 중 하나로 업로드합니다.
-
spark-pipeline.yml및transformations/디렉터리가 포함된.zip파일 -
동일한 구조를 포함하는 Amazon S3 접두사(디렉터리)
4단계: AWS Glue 작업 생성 및 실행
다음 파라미터를 사용하여 AWS Glue 작업을 생성합니다.
-
--enable-spark-declarative-pipeline:true(필수 - SDP 모드를 활성화함) -
ScriptLocation: 파이프라인 정의 zip 또는 Amazon S3 접두사(SDP 파이프라인에 필요) -
--enable-glue-datacatalog:true(선택 사항 - Data Catalog에서 테이블 등록)
다음 예제에서는 AWS CLI를 사용하여 SDP 작업을 생성합니다.
aws glue create-job \ --name my-sdp-pipeline \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 2 \ --command '{"Name":"glueetl","ScriptLocation":"s3://my-bucket/pipelines/my_pipeline.zip"}' \ --default-arguments '{ "--enable-spark-declarative-pipeline": "true", "--enable-glue-datacatalog": "true" }'
파이프라인 실행
StartJobRun을 사용하여 SDP 파이프라인을 실행합니다. 런타임 시 전달된 작업 인수를 사용하여 실행 동작을 제어할 수 있습니다.
실행 모드
파이프라인 실행을 제어하려면 다음 인수를 StartJobRun에 전달합니다.
--conf spark.glue.sdp.jobMode-
실행 모드를 제어합니다.
RUN(기본값) - 파이프라인을 정상적으로 실행합니다.VALIDATE- 데이터를 쓰지 않고 YAML 구문, 종속성 확인 및 SQL/Python 컴파일을 검사하는 모의 실행을 수행합니다.
--conf spark.glue.sdp.runMode-
새로 고치는 데이터세트를 제어합니다.
- 기본값(실행 모드 플래그 없음)
-
모든 데이터세트를 실행합니다. 구체화된 뷰는 전체를 다시 계산하지만, 스트리밍 테이블은 마지막 체크포인트 이후의 새 데이터만 처리합니다.
--refresh <dataset>-
지정된 데이터세트만 새로 고칩니다. 스트리밍 테이블은 새 데이터를 증분 방식으로 처리하고 구체화된 뷰는 전체를 다시 계산합니다.
--full-refresh <dataset>-
지정된 데이터세트만 재설정하고 다시 계산합니다. 스트리밍 테이블의 경우 이 작업은 체크포인트를 재설정하고 모든 데이터를 다시 처리합니다.
--full-refresh-all-
모든 데이터세트를 재설정하고 다시 계산합니다.
SDP에서 Iceberg 테이블 사용
Apache Iceberg는 Hive 또는 AWS Glue 관리형 테이블에서 사용하는 파일 기반 _spark_metadata 로그에 의존하지 않으므로 내구성이 뛰어난 교차 실행 증분 처리가 필요한 스트리밍 테이블에 권장되는 테이블 형식입니다. 출력 테이블을 데이터 카탈로그에 등록할지 여부에 따라 두 가지 방법으로 SDP와 함께 Iceberg를 구성할 수 있습니다.
옵션 1: 데이터 카탈로그와 함께 사용하는 Iceberg(권장)
Iceberg 테이블을 데이터 카탈로그에 등록(table_type=ICEBERG 사용)하여 , Amazon Redshift, Amazon EMR과 같은 다른 엔진에서 쿼리할 수 있도록 하려면 이 옵션을 사용합니다. 테이블 데이터 및 메타데이터는 Amazon S3에 저장되며 교차 실행 증분 상태가 유지됩니다. spark-pipeline.yml의 configuration 섹션에 다음을 추가합니다.
configuration: spark.sql.catalog.glue_catalog: "org.apache.iceberg.spark.SparkCatalog" spark.sql.catalog.glue_catalog.catalog-impl: "org.apache.iceberg.aws.glue.GlueCatalog" spark.sql.catalog.glue_catalog.io-impl: "org.apache.iceberg.aws.s3.S3FileIO" spark.sql.catalog.glue_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
파이프라인 YAML에서 catalog: glue_catalog를 설정하고 database:를 AWS Glue 데이터베이스로 설정합니다. AWS Glue 작업을 생성할 때 --enable-spark-declarative-pipeline를 true으로 설정합니다. Iceberg 테이블에는 --enable-glue-datacatalog를 설정하지 마세요.
참고
Iceberg 카탈로그는 테이블 데이터 및 메타데이터에 자체 warehouse 위치를 사용합니다. 이 옵션을 사용하는 경우 Iceberg 테이블 자체에 대해 spark.sql.warehouse.dir 또는 데이터베이스 LocationUri를 설정할 필요가 없습니다.
옵션 2: 파일 기반(Hadoop) 카탈로그를 사용하는 Iceberg
데이터 카탈로그 등록이 필요하지 않은 경우 이 옵션을 사용합니다. Iceberg 메타데이터는 Amazon S3에서 파일 기반으로 관리되며 테이블은 데이터 카탈로그에 등록되지 않습니다. spark-pipeline.yml의 configuration 섹션에 다음을 추가합니다.
configuration: spark.sql.catalog.spark_catalog: "org.apache.iceberg.spark.SparkSessionCatalog" spark.sql.catalog.spark_catalog.type: "hadoop" spark.sql.catalog.spark_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
참고
이 옵션에 대한 다음 사항을 참고하세요.
-
이 옵션으로 생성된 테이블은 데이터 카탈로그에 등록되지 않으므로와 같은 쿼리 엔진에서 쿼리할 수 없습니다.
-
이 옵션은 테이블 스토리지에 자체
warehouse위치를 사용합니다.
테이블을 데이터 카탈로그에 등록하고 다른 엔진에서 쿼리해야 하는 경우 옵션 1을 사용하세요. 데이터 카탈로그 등록이 필요하지 않은 경우에만 옵션 2를 사용하세요.
참고
--enable-glue-datacatalog 작업 파라미터는 Spark Hive 메타스토어를 Hive(비 Iceberg) 테이블용 데이터 카탈로그에 연결합니다. Iceberg 테이블의 경우 catalog-impl=GlueCatalog는 AWS SDK를 통해 데이터 카탈로그에 직접 테이블을 등록하므로 Iceberg에 대해 --enable-glue-datacatalog를 설정하지 않습니다. Iceberg 테이블을 데이터 카탈로그에 등록하기 위해 기본 SparkSessionCatalog(type: hive)을 --enable-glue-datacatalog와 함께 사용하여 Iceberg를 구성하지 마십시오. AWS Glue 6.0에서는 이 조합이 실패합니다.
Iceberg를 구성하면 다음과 같은 이점을 얻을 수 있습니다.
-
스트리밍 테이블은 작업 실행 간에 Amazon S3에서 체크포인트 상태를 유지합니다.
-
각 실행은 새 데이터 파일과 Iceberg 스냅샷을 생성합니다.
-
후속 실행은 마지막으로 커밋된 오프셋에서 재개됩니다.
-
Iceberg의 스냅샷 메커니즘을 통해 전체 테이블 기록이 보존됩니다.
Iceberg 테이블을 스트리밍 소스로 삼아 증분 방식으로 읽을 수도 있습니다. 메달리온 아키텍처에서 다운스트림 스트리밍 테이블은 매 실행마다 업스트림 Iceberg 테이블에 커밋된 새 행만 사용할 수 있습니다. 다음 예제는 Iceberg bronze 테이블에서 silver 스트리밍 테이블로 증분 방식으로 읽습니다.
from pyspark import pipelines as dp from pyspark.sql import SparkSession spark = SparkSession.active() dp.create_streaming_table( "silver", comment="Incremental silver layer built from the bronze Iceberg table" ) @dp.append_flow(target="silver") def from_bronze(): # Incremental read from the Iceberg bronze table; each run processes only new rows. return spark.readStream.table("bronze")
고려 사항 및 제한
SDP를 사용할 때 다음 사항을 고려하세요.
-
구체화된 뷰는 항상 전체를 다시 계산합니다. 증분 새로 고침은 지원되지 않습니다. 증분 워크로드에는 스트리밍 테이블을 사용하세요.
-
스트리밍 테이블 Python API.
dp.create_streaming_table()와 함께@dp.append_flow(target=...)을 사용합니다. -
스트리밍 테이블에 대한 교차 실행 증분 처리. 스트리밍 테이블은 데이터 및 체크포인트 상태가 Amazon S3에 지속되는 경우에만 교차 실행 증분 처리를 지원합니다. Hive 또는 AWS Glue 관리형(비 Iceberg) 테이블은 데이터베이스
LocationUri를 Amazon S3 경로로 설정해야 하는 반면 Apache Iceberg 테이블은 테이블 메타데이터를 직접 관리합니다. -
Hive 또는 AWS Glue 관리형(비 Iceberg) 테이블에는 Database LocationUri가 필요합니다.
catalog-impl=GlueCatalog와 함께 데이터 카탈로그를 사용하는 Iceberg 테이블에는 필요하지 않습니다. 자세한 내용은 사전 조건을 참조하세요. -
데이터 품질 기대치. 인라인 데이터 품질 주석은 현재 SDP 프레임워크에서 지원되지 않습니다.
-
다운스트림 쿼리 함수에서
withColumn을 사용하지 않음. 다운스트림 데이터세트(예: 구체화된 뷰)가spark.table(...)을 사용하여 업스트림 파이프라인 데이터세트에서 읽고.withColumn(...)을 적용하면 SDP가 두 번째 및 이후의 실행에서 데이터세트 간의 종속성을 감지하지 못할 수 있습니다. 이로 인해 다운스트림이 이전 실행에서 오래된 데이터를 읽습니다(1회 실행 지연). 이 문제를 방지하려면.withColumn(...)을 사용하는 대신.select(...)내부에 파생된 열을 표현하세요. 또한 쿼리 함수 내에서 계획 해석(예:.schema또는.collect)을 강제로 적용하는 작업은 피합니다. -
마이그레이션 도구 없음. 다른 파이프라인 프레임워크에서의 자동 마이그레이션은 지원되지 않습니다. 테이블을 증분식으로 마이그레이션: SDP는 기존 카탈로그 테이블에서 읽을 수 있습니다.
-
일정 예약. SDP 작업은 다른 AWS Glue 작업(AWS Glue 트리거, Amazon EventBridge, Apache Airflow)과 동일한 예약 메커니즘을 사용합니다.
명령형 스크립트에서 SDP로 마이그레이션
기존 명령 Spark 스크립트를 SDP로 점진적으로 마이그레이션할 수 있습니다.
-
단일
spark.sql(...).write.saveAsTable(...)호출을CREATE MATERIALIZED VIEWSQL 문으로 변환하여 테이블 1개로 시작합니다. -
테이블을 증분식으로 추가합니다. SDP는 혼합 종속성을 처리합니다. SDP 테이블은 파이프라인의 일부가 아닌 기존 카탈로그 테이블에서 읽을 수 있습니다.
-
전환 중에 두 패턴을 병렬로 실행합니다. SDP 작업과 명령형 작업은 공존할 수 있습니다.
SDP는 기존 Data Catalog 테이블, 외부 테이블, 데이터베이스 간 참조를 포함하여 SparkSession을 통해 액세스할 수 있는 모든 테이블을 참조할 수 있습니다.