Spark 선언형 파이프라인 프로그래밍 가이드

Spark 선언형 파이프라인 프로그래밍 가이드 (Spark Declarative Pipelines Programming Guide)

Apache Spark에서 신뢰할 수 있고 유지보수하기 쉬우며 테스트 가능한 데이터 파이프라인을 선언적으로(declaratively) 구축하는 프레임워크인 Spark Declarative Pipelines(SDP)에 대한 문서예요. 핵심 개념(Flows, Datasets, Pipelines), spark-pipelines CLI, 그리고 Python·SQL로 파이프라인을 작성하는 방법을 알아볼게요.

출처: 문서

본문

Spark Declarative Pipelines(SDP)이란? (What is Spark Declarative Pipelines?)

Spark Declarative Pipelines(SDP)은 Apache Spark에서 신뢰할 수 있고 유지보수하기 쉬우며 테스트 가능한 데이터 파이프라인을 구축하기 위한 선언형 프레임워크예요. SDP는 파이프라인 실행의 메커니즘보다 데이터에 적용하려는 변환에 집중할 수 있게 해서 ETL 개발을 단순화해요.

SDP는 배치(batch)와 스트리밍 데이터 처리 모두를 위해 설계됐으며, 다음과 같은 일반적인 사용 사례를 지원해요:

  • 클라우드 스토리지에서 데이터 수집 (Amazon S3, Azure ADLS Gen2, Google Cloud Storage)
  • 메시지 버스에서 데이터 수집 (Apache Kafka, Amazon Kinesis, Google Pub/Sub, Azure EventHub)
  • 증분 배치 및 스트리밍 변환

SDP의 핵심 장점은 선언형 접근 방식이에요. 어떤 테이블이 존재해야 하고 그 내용이 무엇이어야 하는지 정의하면, SDP가 오케스트레이션, 컴퓨트 관리, 오류 처리를 자동으로 처리해줘요.

빠른 설치 (Quick install)

SDP를 빠르게 설치하는 방법은 pip을 사용하는 것이에요:

pip install pyspark[pipelines]

더 많은 설치 옵션은 다운로드 페이지를 참고하세요.

핵심 개념 (Key Concepts)

Flows

플로우(flow)는 스트리밍과 배치 의미론(semantics)을 모두 지원하는 SDP의 기초적인 데이터 처리 개념이에요. 플로우는 소스에서 데이터를 읽고, 사용자 정의 처리 로직을 적용한 다음, 결과를 대상 데이터셋에 써요.

예를 들어 다음과 같은 쿼리를 작성하면:

CREATE STREAMING TABLE target_table AS
SELECT * FROM STREAM source_table

SDP는 target_table이라는 테이블을 만들고, source_table에서 새 데이터를 읽어 target_table에 쓰는 플로우를 함께 만들어요.

데이터셋 (Datasets)

데이터셋(dataset)은 파이프라인 내 하나 이상의 플로우의 출력인 쿼리 가능한 객체예요. 파이프라인의 플로우는 파이프라인에서 생성된 데이터셋을 읽을 수도 있어요.

  • 스트리밍 테이블 (Streaming Table) – 테이블의 정의와 그 테이블에 쓰는 하나 이상의 스트리밍 플로우예요. 스트리밍 테이블은 데이터의 증분 처리를 지원하므로, 도착하는 새 데이터만 처리할 수 있어요.
  • 구체화된 뷰 (Materialized View) – 테이블로 사전 계산(precomputed)되는 뷰예요. 구체화된 뷰에는 정확히 하나의 배치 플로우가 항상 쓰여요.
  • 임시 뷰 (Temporary View) – 파이프라인의 한 실행으로 범위가 제한되는 뷰예요. 파이프라인 내의 플로우들에서 참조할 수 있어요. 파이프라인의 다른 여러 요소가 의존하는 변환과 중간 논리 엔터티를 캡슐화하는 데 유용해요.

파이프라인 (Pipelines)

파이프라인(pipeline)은 SDP의 주요 개발·실행 단위예요. 파이프라인은 하나 이상의 플로우, 스트리밍 테이블, 구체화된 뷰를 포함할 수 있어요. 파이프라인이 실행되는 동안 정의된 객체들의 의존성을 분석하고 실행 순서와 병렬화를 자동으로 오케스트레이션해요.

파이프라인 프로젝트 (Pipeline Projects)

파이프라인 프로젝트(pipeline project)는 파이프라인을 구성하는 데이터셋과 플로우의 코드 정의를 담고 있는 소스 파일들의 집합이에요. 소스 파일은 .py 또는 .sql 파일일 수 있어요.

파이프라인 스펙 파일(spec file)의 이름은 spark-pipeline.yml 또는 spark-pipeline.yaml로 짓는 것이 관례예요.

YAML 형식의 파이프라인 스펙 파일은 파이프라인 프로젝트의 최상위 구성을 담고 있으며 다음 필드를 가져요:

  • name (필수) – 파이프라인 프로젝트의 이름이에요.
  • libraries (필수) – 변환 소스 파일(SQL 또는 Python)이 있는 경로들이에요.
  • storage (필수) – 파이프라인 내 스트리밍 테이블에 대한 체크포인트를 저장할 수 있는 디렉터리예요.
  • database (선택) – 파이프라인 출력의 기본 대상 데이터베이스예요. schema를 별칭으로 대신 사용할 수도 있어요.
  • catalog (선택) – 파이프라인 출력의 기본 대상 카탈로그예요.
  • configuration (선택) – Spark 구성 속성의 맵이에요.

파이프라인 스펙 파일의 예:

name: my_pipeline
libraries:
  - glob:
      include: transformations/**
storage: file:///absolute/path/to/storage/dir
catalog: my_catalog
database: my_db
configuration:
  spark.sql.shuffle.partitions: "1000"

아래에서 설명할 spark-pipelines init 명령은 기본 구성과 디렉터리 구조로 파이프라인 프로젝트를 쉽게 생성할 수 있게 해줘요.

spark-pipelines 명령줄 인터페이스 (The spark-pipelines Command Line Interface)

spark-pipelines 명령줄 인터페이스(CLI)는 파이프라인을 관리하는 주요 방법이에요.

spark-pipelinesspark-submit 위에 구축되므로, spark-submit이 지원하는 모든 클러스터 매니저를 지원해요. --class를 제외한 모든 spark-submit 인자를 지원해요.

spark-pipelines init

spark-pipelines init --name my_pipelinemy_pipeline이라는 디렉터리 안에 스펙 파일과 예제 변환 정의를 포함한 간단한 파이프라인 프로젝트를 생성해요.

spark-pipelines run

spark-pipelines run은 파이프라인의 실행을 시작하고 완료될 때까지 그 진행 상황을 모니터링해요.

spark-pipelinesspark-submit 위에 구축되므로 --class를 제외한 모든 spark-submit 인자를 지원해요. 사용 가능한 파라미터의 전체 목록은 Spark Submit 문서를 참고하세요.

또한 다음과 같은 파이프라인 전용 파라미터를 지원해요:

  • --spec PATH – 파이프라인 스펙 파일의 경로예요. 제공되지 않으면 CLI는 현재 디렉터리와 상위 디렉터리에서 다음 파일 중 하나를 찾아요:
    • spark-pipeline.yml
    • spark-pipeline.yaml
  • --full-refresh DATASETS – 재설정·재계산할 데이터셋 목록(쉼표로 구분)이에요. 지정된 데이터셋의 모든 기존 데이터와 체크포인트를 지우고 처음부터 다시 계산해요.
  • --full-refresh-all – 전체 그래프 재설정 및 재계산을 수행해요. 파이프라인의 모든 데이터셋에 --full-refresh를 적용하는 것과 동일해요.
  • --refresh DATASETS – 갱신할 데이터셋 목록(쉼표로 구분)이에요. 기존 데이터를 지우지 않고 지정된 데이터셋에 대한 갱신을 트리거해요.
갱신 선택 동작 (Refresh Selection Behavior)

갱신 옵션을 지정하지 않으면 기본적으로 증분 갱신(incremental update)이 수행돼요. 갱신 파라미터는 상호 배타적이에요:

  • --full-refresh-all--full-refresh 또는 --refresh와 함께 사용할 수 없어요.
  • --full-refresh--refresh는 서로 다른 데이터셋에 서로 다른 동작을 지정하기 위해 함께 사용할 수 있어요.
예제 (Examples)
# Basic run with default incremental update
spark-pipelines run

# Run with specific spec file
spark-pipelines run --spec /path/to/my-pipeline.yaml

# Full refresh of specific datasets
spark-pipelines run --full-refresh orders,customers

# Full refresh of entire pipeline
spark-pipelines run --full-refresh-all

# Run with custom Spark configuration
spark-pipelines run --conf spark.sql.shuffle.partitions=200 --driver-memory 4g

# Run on remote Spark Connect server
spark-pipelines run --remote sc://my-cluster:15002

spark-pipelines dry-run

spark-pipelines dry-run은 데이터를 읽거나 쓰지 않지만, 파이프라인이 실제로 실행됐다면 잡혔을 많은 종류의 오류를 잡아내는 파이프라인 실행을 시작해요. 예를 들어:

  • 구문 오류 – 예: 잘못된 Python 또는 SQL 코드
  • 분석 오류 – 예: 존재하지 않는 테이블이나 컬럼에서 select
  • 그래프 검증 오류 – 예: 순환 의존성(cyclic dependencies)

spark-pipelinesspark-submit 위에 구축되므로 --class를 제외한 모든 spark-submit 인자를 지원해요. 사용 가능한 파라미터의 전체 목록은 Spark Submit 문서를 참고하세요.

또한 파이프라인 전용 --spec 파라미터를 지원해요 (설명은 위 run 섹션 참고).

Python으로 SDP 프로그래밍 (Programming with SDP in Python)

SDP Python 정의는 pyspark.pipelines 모듈에 정의돼 있어요.

Python API로 구현한 파이프라인은 이 모듈을 import해야 해요. 모듈을 dp로 별칭 지을 것을 권장해요.

from pyspark import pipelines as dp

Python 파이프라인의 Spark 세션 (The Spark Session in Python Pipelines)

Spark 4.1에서는 모든 파이프라인 파일이 spark = SparkSession.active()를 명시적으로 선언해야 했어요. Spark 4.2부터 프레임워크가 각 파이프라인 파일의 모듈 네임스페이스에 spark를 주입하므로, 명시적 할당이 더 이상 필요하지 않아요.

from pyspark import pipelines as dp

@dp.materialized_view
def my_view():
    return spark.range(10)

여전히 spark = SparkSession.active()를 포함하는 파이프라인 파일은 올바르게 동작해요. 하지만 세션을 명시적으로 할당한다면 SparkSession.active()가 유일하게 지원되는 방법이에요. 예를 들어 SparkSession.builder.config(...).getOrCreate()는 세션 구성을 변경하는데, 이는 SDP에서 차단돼요.

명시적 할당이 없으면 많은 도구와 편집기가 spark를 정의되지 않은 이름으로 간주할 수 있어요. 이를 해결하려면 모듈 스코프에 spark: SparkSession을 추가하면 돼요. SDP는 모듈이 실행되기 전에 실제 세션을 여전히 주입하므로, 이는 정적 분석을 위한 타입 문서화일 뿐이에요.

from pyspark import pipelines as dp
from pyspark.sql import SparkSession

spark: SparkSession

@dp.materialized_view
def my_view():
    return spark.range(10)

Python에서 구체화된 뷰 만들기 (Creating a Materialized View in Python)

@dp.materialized_view 데코레이터는 배치 읽기를 수행하는 함수의 결과에 기반해 구체화된 뷰를 생성하도록 SDP에 알려줘요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.materialized_view
def basic_mv() -> DataFrame:
    return spark.table("samples.nyctaxi.trips")

구체화된 뷰의 이름은 함수의 이름에서 파생돼요.

name 인자로 구체화된 뷰의 이름을 지정할 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.materialized_view(name="trips_mv")
def basic_mv() -> DataFrame:
    return spark.table("samples.nyctaxi.trips")

Python에서 임시 뷰 만들기 (Creating a Temporary View in Python)

@dp.temporary_view 데코레이터는 배치 읽기를 수행하는 함수의 결과에 기반해 임시 뷰를 생성하도록 SDP에 알려줘요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.temporary_view
def basic_tv() -> DataFrame:
    return spark.table("samples.nyctaxi.trips")

이 임시 뷰는 파이프라인 내의 다른 쿼리에서 읽을 수 있지만, 파이프라인의 범위 밖에서는 읽을 수 없어요.

Python에서 스트리밍 테이블 만들기 (Creating a Streaming Table in Python)

스트리밍 읽기를 수행하는 함수와 함께 @dp.table 데코레이터를 사용해 스트리밍 테이블을 만들 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.table
def basic_st() -> DataFrame:
    return spark.readStream.table("samples.nyctaxi.trips")

Python에서 스트리밍 소스에서 데이터 로드하기 (Loading Data from Streaming Sources in Python)

SDP는 Spark Structured Streaming(spark.readStream)이 지원하는 모든 형식에서 데이터 로드를 지원해요.

예를 들어 Kafka 토픽에서 읽는 쿼리를 가진 스트리밍 테이블을 만들 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.table
def ingestion_st() -> DataFrame:
    return (
        spark.readStream
        .format("kafka")
        .option("kafka.bootstrap.servers", "localhost:9092")
        .option("subscribe", "orders")
        .load()
    )

Python에서 배치 소스에서 데이터 로드하기 (Loading Data from Batch Sources in Python)

SDP는 Spark SQL(spark.read)이 지원하는 모든 형식에서 데이터 로드를 지원해요.

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

@dp.materialized_view
def batch_mv() -> DataFrame:
    return spark.read.format("json").load("/datasets/retail-org/sales_orders")

Python에서 파이프라인에 정의된 테이블 조회하기 (Querying Tables Defined in a Pipeline in Python)

파이프라인에 정의된 다른 테이블을, 파이프라인 밖에서 정의된 테이블을 참조하는 것과 같은 방식으로 참조할 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame
from pyspark.sql.functions import col

@dp.table
def orders() -> DataFrame:
    return (
        spark.readStream
        .format("kafka")
        .option("kafka.bootstrap.servers", "localhost:9092")
        .option("subscribe", "orders")
        .load()
    )

@dp.materialized_view
def customers() -> DataFrame:
    return (
        spark.read
        .format("csv")
        .option("header", True)
        .load("/datasets/retail-org/customers")
    )

@dp.materialized_view
def customer_orders() -> DataFrame:
    return (
        spark.table("orders")
        .join(
            spark.table("customers"), "customer_id")
            .select(
                "customer_id",
                "order_number",
                "state",
                col("order_datetime").cast("date").alias("order_date"),
            )
        )
    )

@dp.materialized_view
def daily_orders_by_state() -> DataFrame:
    return (
        spark.table("customer_orders")
        .groupBy("state", "order_date")
        .count()
        .withColumnRenamed("count", "order_count")
    )

Python에서 for 루프로 테이블 만들기 (Creating Tables in For Loop in Python)

Python for 루프를 사용해 프로그래밍 방식으로 여러 테이블을 만들 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame
from pyspark.sql.functions import collect_list, col

@dp.temporary_view()
def customer_orders() -> DataFrame:
    orders = spark.table("samples.tpch.orders")
    customer = spark.table("samples.tpch.customer")

    return (
        orders
        .join(customer, orders.o_custkey == customer.c_custkey)
        .select(
            col("c_custkey").alias("custkey"),
            col("c_name").alias("name"),
            col("c_nationkey").alias("nationkey"),
            col("c_phone").alias("phone"),
            col("o_orderkey").alias("orderkey"),
            col("o_orderstatus").alias("orderstatus"),
            col("o_totalprice").alias("totalprice"),
            col("o_orderdate").alias("orderdate"),
        )
    )

@dp.temporary_view()
def nation_region() -> DataFrame:
    nation = spark.table("samples.tpch.nation")
    region = spark.table("samples.tpch.region")

    return (
        nation
        .join(region, nation.n_regionkey == region.r_regionkey)
        .select(
            col("n_name").alias("nation"),
            col("r_name").alias("region"),
            col("n_nationkey").alias("nationkey"),
        )
    )

# Extract region names from region table
region_list = spark.table("samples.tpch.region").select(collect_list("r_name")).collect()[0][0]

# Iterate through region names to create new region-specific materialized views
for region in region_list:
    @dp.table(name=f"{region.lower().replace(' ', '_')}_customer_orders")
    def regional_customer_orders(region_filter=region) -> DataFrame:
        customer_orders = spark.table("customer_orders")
        nation_region = spark.table("nation_region")

        return (
            customer_orders
            .join(nation_region, customer_orders.nationkey == nation_region.nationkey)
            .select(
                col("custkey"),
                col("name"),
                col("phone"),
                col("nation"),
                col("region"),
                col("orderkey"),
                col("orderstatus"),
                col("totalprice"),
                col("orderdate"),
            )
            .filter(f"region = '{region_filter}'")
        )

Python에서 단일 대상에 여러 플로우 쓰기 (Using Multiple Flows to Write to a Single Target in Python)

같은 데이터셋에 데이터를 추가(append)하는 여러 플로우를 만들 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql import DataFrame

# create a streaming table
dp.create_streaming_table("customers_us")

# define the first append flow
@dp.append_flow(target = "customers_us")
def append_customers_us_west() -> DataFrame:
    return spark.readStream.table("customers_us_west")

# define the second append flow
@dp.append_flow(target = "customers_us")
def append_customers_us_east() -> DataFrame:
    return spark.readStream.table("customers_us_east")

SQL로 SDP 프로그래밍 (Programming with SDP in SQL)

SQL에서 구체화된 뷰 만들기 (Creating a Materialized View in SQL)

SQL로 구체화된 뷰를 만드는 기본 구문은 다음과 같아요:

CREATE MATERIALIZED VIEW basic_mv
AS SELECT * FROM samples.nyctaxi.trips;

SQL에서 임시 뷰 만들기 (Creating a Temporary View in SQL)

SQL로 임시 뷰를 만드는 기본 구문은 다음과 같아요:

CREATE TEMPORARY VIEW basic_tv
AS SELECT * FROM samples.nyctaxi.trips;

SQL에서 스트리밍 테이블 만들기 (Creating a Streaming Table in SQL)

스트리밍 테이블을 만들 때는 STREAM 키워드를 사용해 소스에 대한 스트리밍 의미론을 나타내요:

CREATE STREAMING TABLE basic_st
AS SELECT * FROM STREAM samples.nyctaxi.trips;

SQL에서 파이프라인에 정의된 테이블 조회하기 (Querying Tables Defined in a Pipeline in SQL)

파이프라인에 정의된 다른 테이블을 참조할 수 있어요:

CREATE STREAMING TABLE orders
AS SELECT * FROM STREAM orders_source;

CREATE MATERIALIZED VIEW customers
AS SELECT * FROM customers_source;

CREATE MATERIALIZED VIEW customer_orders
AS SELECT
  c.customer_id,
  o.order_number,
  c.state,
  date(timestamp(int(o.order_datetime))) order_date
FROM orders o
INNER JOIN customers c
ON o.customer_id = c.customer_id;

CREATE MATERIALIZED VIEW daily_orders_by_state
AS SELECT state, order_date, count(*) order_count
FROM customer_orders
GROUP BY state, order_date;

SQL에서 단일 대상에 여러 플로우 쓰기 (Using Multiple Flows to Write to a Single Target in SQL)

같은 대상에 데이터를 추가하는 여러 플로우를 만들 수 있어요:

-- create a streaming table
CREATE STREAMING TABLE customers_us;

-- define the first append flow
CREATE FLOW append_customers_us_west
AS INSERT INTO customers_us
SELECT * FROM STREAM(customers_us_west);

-- define the second append flow
CREATE FLOW append_customers_us_east
AS INSERT INTO customers_us
SELECT * FROM STREAM(customers_us_east);

싱크로 외부 대상에 데이터 쓰기 (Writing Data to External Targets with Sinks)

SDP의 싱크(Sink)는 기본 스트리밍 테이블과 구체화된 뷰 너머의 외부 대상에 변환된 데이터를 쓰는 방법을 제공해요. 싱크는 낮은 지연시간 데이터 처리를 요구하는 운영(operational) 사용 사례, reverse ETL 연산, 또는 외부 시스템에 쓰기에 특히 유용해요.

싱크는 Spark Structured Streaming 쿼리를 쓸 수 있는 모든 대상에 파이프라인이 쓸 수 있게 해주며, 여기에는 Apache Kafka와 Azure Event Hubs가 포함되지만 이에 국한되지 않아요.

Python에서 싱크 만들고 사용하기 (Creating and Using Sinks in Python)

싱크를 사용하려면 두 가지 주요 단계가 있어요: 싱크 정의를 만들고, 데이터를 쓰기 위한 append 플로우를 구현하는 것이에요.

Kafka 싱크 만들기 (Creating a Kafka Sink)

Kafka 토픽으로 데이터를 스트리밍하는 싱크를 만들 수 있어요:

from pyspark import pipelines as dp
from pyspark.sql.functions import to_json, struct

dp.create_sink(
    name="kafka_sink",
    format="kafka",
    options={
        "kafka.bootstrap.servers": "localhost:9092",
        "topic": "processed_orders"
    }
)

@dp.append_flow(target="kafka_sink")
def kafka_orders_flow() -> DataFrame:
    return (
        spark.readStream.table("customer_orders")
        .select(
            col("order_id").cast("string").alias("key"),
            to_json(struct("*")).alias("value")
        )
    )

싱크 고려 사항 (Sink Considerations)

싱크를 사용할 때 다음 고려 사항을 기억하세요:

  • 스트리밍 전용: 싱크는 현재 append_flow 데코레이터를 통한 스트리밍 쿼리만 지원해요.
  • Python API: 싱크 기능은 SQL이 아닌 Python API에서만 사용할 수 있어요.
  • 추가 전용: append 연산만 지원돼요. full refresh 갱신은 체크포인트를 재설정하지만, 이전에 계산된 결과는 정리하지 않아요.

중요한 고려 사항 (Important Considerations)

Python 고려 사항 (Python Considerations)

  • SDP는 파이프라인을 정의하는 코드를 계획(planning)과 파이프라인 실행 중에 여러 번 평가해요. 데이터셋을 정의하는 Python 함수는 테이블이나 뷰를 정의하는 데 필요한 코드만 포함해야 해요.
  • 데이터셋을 정의하는 데 사용되는 함수는 pyspark.sql.DataFrame을 반환해야 해요.
  • SDP 데이터셋 코드의 일부로 파일이나 테이블에 저장하거나 쓰는 메서드는 절대 사용하지 마세요.
  • Python에서 데이터셋을 정의하기 위해 for 루프 패턴을 사용할 때는 for 루프에 전달되는 값 목록이 항상 가산적(additive)인지 확인하세요.

SDP 코드에서 절대 사용하면 안 되는 Spark SQL 연산의 예:

  • collect()
  • count()
  • pivot()
  • toPandas()
  • save()
  • saveAsTable()
  • start()
  • toTable()

SQL 고려 사항 (SQL Considerations)

  • PIVOT 절은 SDP SQL에서 지원되지 않아요.

더 알아보기 (Learn more)