독립 실행형 파이프라인에서 Python 사용

Python 사용하여 Notebook에서 독립 실행형 구체화된 뷰 및 스트리밍 테이블을 만들고 새로 고칠 수 있습니다. 이렇게 하면 다른 Python 기반 Notebook 워크플로와 함께 독립 실행형 파이프라인을 관리할 수 있습니다.

이 작업을 수행하는 방법에는 두 가지가 있습니다.

  • pyspark.pipelines, @dp.materialized_view 및 @dp.table 데코레이터를 사용하여 테이블을 정의합니다. 이 로직이 DataFrame 코드로 표현하기 더 쉬울 때 사용하세요. 파이프라인 데코레이터를 사용하여 테이블 정의하기를 참조하세요.
  • Databricks SQL 웨어하우스에서 실행하는 것과 동일한 SQL 문을 spark.sql()에 전달하여 전송하세요. 이렇게 하면 REFRESH 구문과 새로고침 일정을 포함한 독립 실행형 구체화된 뷰 및 스트리밍 테이블의 전체 SQL 인터페이스를 확인할 수 있습니다. spark.sql()를 사용하여 SQL 문 제출을 참조하세요.

독립 실행형 파이프라인용 Python 소스에는 서버리스 일반 컴퓨팅에 연결된 노트북이 필요합니다. Python 사용하여 Databricks SQL 웨어하우스에서 독립 실행형 파이프라인을 만들거나 새로 고칠 수 없습니다. 웨어하우스는 Python Notebook이 아닌 SQL 문을 실행하기 때문입니다. 대신 SQL 웨어하우스를 사용하려면 독립 실행형 구체화된 뷰 사용 및 독립 실행형 스트리밍 테이블 사용을 참조하세요.

Important

서버리스 일반 컴퓨팅의 Notebook에서 독립 실행형 구체화된 뷰 및 스트리밍 테이블을 만들고 새로 고치는 기능은 베타 로 제공되며 일부 지역에서 사용할 수 있습니다. 노트북을 참조하세요.

Requirements

Python 사용하여 독립 실행형 파이프라인을 만들고 새로 고치려면 Databricks Runtime 18.1 이상의 서버리스 일반 컴퓨팅에 연결된 Notebook이 필요합니다. 지역별 가용성 및 사용 권한을 포함한 전체 요구 사항 목록은 Notebooks를 참조하세요.

파이프라인 데코레이터가 포함된 테이블을 정의하세요

Lakeflow 파이프라인에서 사용하는 데코레이터와 동일한 도구로 독립형 물질화된 뷰나 스트리밍 테이블을 정의할 수 있습니다. 각 장식된 함수는 하나의 테이블을 정의합니다. 셀을 실행하면 Azure Databricks가 테이블을 생성하고 서버리스 파이프라인을 실행해 이를 채웁니다. 업데이트가 완료되면 셀이 다시 표시됩니다.

Warning

파이프라인 데코레이터는 서버 없는 환경 버전 5 이상을 요구합니다.

구체화된 뷰를 정의하세요

배치 DataFrame을 반환하는 함수에 @dp.materialized_view를 사용하세요. 다음 예시는 Wanderbricks 샘플 데이터셋의 표에서 bookings 물질화된 뷰 daily_booking_revenue 를 생성합니다:

from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.materialized_view(name="main.default.daily_booking_revenue")
def daily_booking_revenue():
  return (
    spark.read.table("samples.wanderbricks.bookings")
    .groupBy("check_in")
    .agg(F.sum("total_amount").alias("total_revenue"))
  )

스트리밍 리드에서 테이블을 정의하려면 대신 사용 @dp.table 하세요.

스트리밍 테이블을 정의하세요

스트리밍 DataFrame을 반환하는 함수에 @dp.table를 사용하세요. 다음 예시는 같은 bookings 테이블의 스트리밍 읽기에서 스트리밍 테이블 bookings_raw 을 생성합니다:

from pyspark import pipelines as dp

@dp.table(name="main.default.bookings_raw")
def bookings_raw():
  return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

함수가 배치 데이터프레임 @dp.table 을 반환하면 대신 물질화된 뷰를 생성합니다. 단 한 가지 예외는 replace_where, 항상 스트리밍 테이블로 이어집니다. 다음 예시는 2025년 7월 1일 이후 체크인의 일일 수익을 이전 날짜를 재계산하지 않고 최신 상태로 유지합니다:

from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.table(
  name="main.default.booking_revenue_rw",
  replace_where=F.col("check_in") >= F.to_date(F.lit("2025-07-01")),
)
def booking_revenue_rw():
  return (
    spark.read.table("samples.wanderbricks.bookings")
    .groupBy("check_in")
    .agg(F.sum("total_amount").alias("total_revenue"))
  )

각 실행은 술어와 일치하는 행을 삭제하고 그 범위만 재계산합니다. REPLACE WHERE 흐름을 사용한 Batch 처리를 참조하세요.

테이블 새로고침하기

데코레이터로 정의한 테이블을 새로고침하려면, 예를 들어 노트북 셀을 다시 실행하거나, 노트북 전체를 실행하거나, 노트북을 작업으로 실행하는 등 그 테이블을 다시 실행하세요. 각 실행은 테이블이 존재하지 않으면 생성하고, 존재하면 새로고침합니다.

소스에서 사용 가능한 모든 데이터를 재처리하려면 두 데코레이터 중 하나에 full_refresh=True를 전달하세요:

@dp.table(name="main.default.bookings_raw", full_refresh=True)
def bookings_raw():
  return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

데코레이터로 정의된 테이블에서는 REFRESH 문을 사용할 수 없고, SCHEDULE 또는 TRIGGER ON UPDATE로 새로 고침을 예약할 수도 없습니다. 일정을 새로고침하려면 SQL로 테이블을 정의하거나 노트북을 작업으로 예약하세요. Lakeflow 작업을 참조하세요.

테이블 구성

데코레이터는 파이프라인 내에서 받는 것과 동일한 공통 데이터셋 매개변수를 받으며, 여기에는 comment, table_properties, partition_cols, cluster_by, schema, 그리고 spark_conf가 포함됩니다:

@dp.materialized_view(
  name="main.default.daily_booking_revenue",
  comment="Daily booking revenue.",
  table_properties={"quality": "gold"},
  cluster_by=["check_in"],
)
def daily_booking_revenue():
  return (
    spark.read.table("samples.wanderbricks.bookings")
    .groupBy("check_in")
    .agg(F.sum("total_amount").alias("total_revenue"))
  )

매개변수 목록은 materialized_view 와 표를 참조하세요.

private=True 지원되지 않는데, 이는 프라이빗 테이블이 같은 파이프라인 내 다른 데이터셋에서만 읽을 수 있기 때문입니다.

지원되지 않는 API

독립형 테이블은 단일 흐름을 가진 단일 데이터셋이기 때문에, 데이터셋 간 관계를 설명하는 API는 제공되지 않습니다. 다음은 파이프라인 외부에서 오류를 발생시키는 경우입니다:

  • @dp.temporary_view 및 dp.create_streaming_table
  • @dp.append_flow 및 기타 추가 흐름
  • dp.create_auto_cdc_flow 및 dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow 그리고 replace_using REPLACE USING flows를 정의하는 매개변수입니다. REPLACE USING 흐름을 사용한 부분 스냅샷 교체를 참조하세요.
  • dp.create_sink
  • 기대치, 예를 들어 @dp.expect 와 @dp.expect_or_fail

이 자료를 사용하려면 Lakeflow 파이프라인을 작성하세요. Python을 사용하여 파이프라인 코드 개발을 참조하세요.

SQL 문을 제출하세요 spark.sql()

Python 노트북에서 Databricks SQL 웨어하우스에서 실행하는 것과 동일한 SQL 구문을 spark.sql()에 전달하세요. 독립 실행형 구체화된 뷰 및 스트리밍 테이블 구문은 동일합니다. 문을 제출하는 방법만 다릅니다. 웨어하우스와 마찬가지로 각 CREATE 또는 REFRESH 문은 서버를 사용하지 않는 파이프라인을 실행하여 작업을 처리합니다.

세션은 spark 기본적으로 Azure Databricks Notebook에서 사용할 수 있으므로 가져올 필요가 없습니다.

구체화된 뷰 만들기

다음 예제에서는 기본 테이블 mv1구체화된 뷰 base_table1 만듭니다.

spark.sql("""
  CREATE OR REPLACE MATERIALIZED VIEW mv1
  AS SELECT
    date,
    sum(sales) AS sum_of_sales
  FROM base_table1
  GROUP BY date
""")

예약된 새로 고침 및 트리거된 새로 고침과 같은 자세한 CREATE MATERIALIZED VIEW 내용은 구체화된 뷰 만들기를 참조하세요.

스트리밍 테이블 만들기

다음 예제에서는 sales 테이블에서 스트리밍 테이블 raw_data을 생성합니다.

spark.sql("""
  CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT product, price FROM STREAM raw_data
""")

자동 로더를 사용하여 파일 로드 및 예약을 비롯한 자세한 CREATE STREAMING TABLE 내용은 독립 실행형 스트리밍 테이블 사용을 참조하세요.

구체화된 뷰 또는 스트리밍 테이블 새로 고침

REFRESH 문을 사용하여 독립 실행형 테이블을 원본의 최신 데이터로 업데이트합니다.

spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

서버리스 일반 컴퓨팅에서는 새로 고침이 동기적입니다. 비동기 새로 고침( ASYNC 키워드)은 지원되지 않습니다. 서버리스 일반 컴퓨팅을 참조하세요.

매개 변수화된 문

Python 코드의 값을 하드 코딩하는 대신 문으로 전달하려면 SQL에서 명명된 매개 변수 표식을 사용하고 인수를 통해 argsspark.sql()해당 값을 제공합니다. 리터럴 값에는 :min_sales와 같은 표식을 직접 사용하세요. 식별자는 일반 문자열 값으로 대체할 수 없으므로, 매개 변수가 테이블, 뷰 또는 스키마와 같은 객체 이름인 경우에만 마커를 IDENTIFIER()로 감싸세요.

다음 예제에서는 구체화된 뷰 이름과 필터 값을 모두 매개 변수화합니다.

mv_name = "main.sales.regional_sales"
min_sales = 1000

spark.sql("""
  CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
  AS SELECT
    region,
    sum(sales) AS sum_of_sales
  FROM base_table1
  WHERE sales > :min_sales
  GROUP BY region
""", args={
  "mv": mv_name,
  "min_sales": min_sales,
})

자세한 내용은 매개 변수 표식 및 IDENTIFIER 절을 참조하세요.

다른 명령문 실행

Python 노트북에서 이를 spark.sql()에 전달하여 새로 고침 예약, 테이블 변경 또는 테이블 삭제 문을 포함한 독립 실행형 구체화된 뷰 또는 스트리밍 테이블 문을 실행할 수 있습니다. SQL 구문을 포함하여 구체화된 뷰 및 스트리밍 테이블을 사용하는 방법을 이해하려면 독립 실행형 구체화된 뷰 사용 및 독립 실행형 스트리밍 테이블 사용을 참조하세요.

Limitations

서버리스 일반 컴퓨팅에서 만든 독립 실행형 구체화된 뷰 및 스트리밍 테이블에는 비동기 새로 고침을 지원하지 않고 테이블당 비용 특성이 없는 등의 추가 제한 사항이 있습니다. 전체 목록은 서버리스 일반 컴퓨팅을 참조하세요.

이 파이프라인들은 SQL 웨어하우스가 아닌 서버리스 일반 컴퓨팅에서 실행되기 때문에, 인클로잉 웨어하우스에서 사용자 지정 태그를 상속받지 않습니다. system.billing.usage에 대한 웨어하우스 태그 전파는 SQL 웨어하우스에서 실행되는 구체화된 뷰와 스트리밍 테이블에만 적용됩니다. 사용자 지정 태그로 SQL 웨어하우스 비용 할당하기를 참조하세요.

파이프라인 데코레이터로 정의된 테이블에는 다음과 같은 추가 제한이 있습니다:

  • REFRESH 문으로는 이를 새로 고칠 수 없으며, SCHEDULE 또는 TRIGGER ON UPDATE로 새로 고침을 예약할 수도 없습니다. 테이블 새로 고침을 참조하세요.
  • 기대치, 추가 플로우, 변경 데이터 캡처(CDC) 플로우, 싱크, 임시 뷰는 지원되지 않습니다. 지원되지 않는 API를 참조하세요.
  • private=True은 지원되지 않습니다.

추가 리소스