Microsoft Fabric Spark 설정

Microsoft Fabric Spark 설정

dbt-fabricspark 어댑터에서 table 구성과 incremental 모델 전략(append, insert_overwrite, merge, microbatch)을 다루는 페이지예요. 특히 Delta 파일 형식과 파티셔닝 설정이 중요해요.

출처: 문서

본문

테이블 구성 (Configuring tables)

모델을 table로 materialize할 때, 표준 모델 config 외에도 dbt-spark 플러그인 특유의 여러 선택적 config를 포함할 수 있어요.

옵션 설명 필수 여부? 예시
file_format 테이블을 만들 때 사용할 파일 형식(parquet, delta, csv)이에요. 선택 delta
location_root 1 테이블 데이터를 저장하는 데 사용하는 지정 디렉터리예요. 테이블 별칭이 여기에 추가돼요. 선택 Files/<folder> 또는 Tables/<tableName>
partition_by 지정한 컬럼으로 테이블을 파티셔닝해요. 각 파티션마다 디렉터리가 만들어져요. 선택 date_day
clustered_by 테이블의 각 파티션을 지정한 컬럼으로 고정된 수의 버킷으로 나눠요. 선택 country_code
buckets 클러스터링 중 생성할 버킷 수예요. clustered_by 지정 시 필수 8
tblproperties 테이블 동작을 구성하는 테이블 속성이에요. 속성은 파일 형식에 따라 다르며, 참조 문서(Parquet, Delta)를 참고하세요. 선택 Provider=delta Location=abfss://.../Files/tables/sales_data TableProperty.created.by=data_engineering_team TableProperty.purpose=sales analytics CreatedBy=Delta Lake CreatedAt=2024-12-01 14:21:00 Format=Parquet PartitionColumns=region MinReaderVersion=1 MinWriterVersion=2

Incremental 모델 (Incremental models)

dbt는 내장 구성과 materialization을 통해 유용하고 직관적인 모델링 추상화를 제공하려 해요. 세상에 존재하는 Spark 클러스터 간 편차가 크고—Delta 파일 형식과 커스텀 런타임이 오픈 소스 사용자에게 제공하는 강력한 기능도 언급할 필요가 있죠—모든 옵션을 이해하는 것 자체가 하나의 과업이에요.

그래서 dbt-fabricspark 플러그인은 incremental_strategy config를 많이 활용해요. 이 config는 incremental materialization이 첫 실행 이후의 실행에서 모델을 어떻게 빌드할지 알려줘요. 다음 세 값 중 하나로 설정할 수 있어요:

  • append (기본값): 기존 데이터를 업데이트하거나 덮어쓰지 않고 새 레코드를 삽입해요.
  • insert_overwrite: partition_by가 지정되면 테이블의 파티션을 새 데이터로 덮어써요. partition_by가 없으면 테이블 전체를 새 데이터로 덮어써요.
  • merge (Delta 파일 형식 전용): unique_key를 기준으로 레코드를 매칭해요. 이전 레코드를 업데이트하고 새 레코드를 삽입해요. (unique_key를 지정하지 않으면 append와 비슷하게 모든 새 데이터를 삽입해요.)
  • microbatch: event_time을 사용해 데이터를 필터링할 시간 기반 범위를 정의하는 microbatch 전략을 구현해요.

이 각각의 전략에는 장단점이 있으며 아래에서 논의할게요. 다른 모델 config와 마찬가지로, incremental_strategydbt_project.yml 또는 모델 파일의 config() 블록 안에서 지정할 수 있어요.

append 전략

append 전략을 따르면 dbt는 모든 새 데이터로 insert into 문을 실행해요. 이 전략의 장점은 모든 플랫폼, 파일 유형, 연결 방법, Fabric Spark 버전에서 단순하고 동작한다는 거예요. 하지만 이 전략은 기존 데이터를 업데이트, 덮어쓰기, 삭제할 수 없어서, 많은 데이터 소스에서 중복 레코드를 삽입할 가능성이 있어요.

append를 incremental 전략으로 지정하는 것은 선택 사항이에요. 지정하지 않았을 때 사용되는 기본 전략이니까요.

소스 코드

fabricspark_incremental.sql

{{ config(
    materialized='incremental',
    incremental_strategy='append',
) }}

--  All rows returned by this query will be appended to the existing table

select * from {{ ref('events') }}
{% if is_incremental() %}
  where event_ts > (select max(event_ts) from {{ this }})
{% endif %}

실행 코드

fabricspark_incremental.sql

create temporary view fabricspark_incremental__dbt_tmp as

    select * from analytics.events

    where event_ts >= (select max(event_ts) from {{ this }})

;

insert into table analytics.fabricspark_incremental
    select `date_day`, `users` from spark_incremental__dbt_tmp

insert_overwrite 전략

이 전략은 모델 config에 partition_by 절과 함께 지정할 때 가장 효과적이에요. dbt는 쿼리에 포함된 모든 파티션을 동적으로 교체하는 원자적 insert overwrite을 실행해요. 이 incremental 전략을 사용할 때는 파티션의 모든 관련 데이터를 다시 선택해야 해요.

partition_by를 지정하지 않으면 insert_overwrite 전략은 테이블의 모든 내용을 원자적으로 교체해, 기존 데이터를 새 레코드로만 덮어써요. 하지만 테이블의 컬럼 스키마는 그대로 유지돼요. 이는 테이블 내용이 덮어써지는 동안 다운타임을 최소화하므로 일부 제한된 상황에서 바람직할 수 있어요. 다른 데이터베이스에서 truncate + insert를 실행하는 것과 비슷한 연산이에요. Delta 형식 테이블의 원자적 교체에는 table materialization(create or replace 실행)을 사용하세요.

사용 참고:

  • 이 전략은 file_format: delta 테이블에서는 지원되지 않아요.

소스 코드

fabricspark_incremental.sql

{{ config(
    materialized='incremental',
    partition_by=['date_day'],
    file_format='parquet'
) }}

/*
  Every partition returned by this query will be overwritten
  when this model runs
*/

with new_events as (

    select * from {{ ref('events') }}

    {% if is_incremental() %}
    where date_day >= date_add(current_date, -1)
    {% endif %}

)

select
    date_day,
    count(*) as users

from events
group by 1

실행 코드

fabricspark_incremental.sql

create temporary view fabricspark_incremental__dbt_tmp as

    with new_events as (

        select * from analytics.events


        where date_day >= date_add(current_date, -1)


    )
    select
        date_day,
        count(*) as users

    from events
    group by 1

;

insert overwrite table analytics.fabricspark_incremental
    partition (date_day)
    select `date_day`, `users` from spark_incremental__dbt_tmp

merge 전략

사용 참고: merge incremental 전략에는 다음이 필요해요:

  • file_format: delta
  • delta 파일 형식용 Fabric Spark Runtime 3.0 이상

dbt는 Fabric Warehouse나 SQL 데이터베이스, Snowflake, BigQuery의 기본 merge 동작과 거의 동일하게 보이는 원자적 merge 문을 실행해요. unique_key가 지정되면(권장), dbt는 키 컬럼에서 매칭되는 새 레코드의 값으로 이전 레코드를 업데이트해요. unique_key를 지정하지 않으면 dbt는 매칭 기준을 생략하고 단순히 모든 새 레코드를 삽입해요(append 전략과 비슷).

소스 코드

merge_incremental.sql

{{ config(
    materialized='incremental',
    file_format='delta',
    unique_key='user_id',
    incremental_strategy='merge'
) }}

with new_events as (

    select * from {{ ref('events') }}

    {% if is_incremental() %}
    where date_day >= date_add(current_date, -1)
    {% endif %}

)

select
    user_id,
    max(date_day) as last_seen

from events
group by 1

실행 코드

target/run/merge_incremental.sql

create temporary view merge_incremental__dbt_tmp as

    with new_events as (

        select * from analytics.events


        where date_day >= date_add(current_date, -1)


    )
    select
        user_id,
        max(date_day) as last_seen

    from events
    group by 1

;

merge into analytics.merge_incremental as DBT_INTERNAL_DEST
    using merge_incremental__dbt_tmp as DBT_INTERNAL_SOURCE
    on DBT_INTERNAL_SOURCE.user_id = DBT_INTERNAL_DEST.user_id
    when matched then update set *
    when not matched then insert *

모델 설명 영속화 (Persisting model descriptions)

관계 레벨 docs 영속화는 dbt에서 지원돼요. docs 영속화 구성에 대한 자세한 내용은 the docs를 참고하세요.

persist_docs 옵션을 적절히 구성하면, describe [table] extended 또는 show table extended in [database] like '*'Comment 필드에서 모델 설명을 볼 수 있어요.

항상 schema, 절대 database 아님

Fabric Spark는 "schema"와 "database"라는 용어를 같은 뜻으로 사용해요. dbt는 databaseschema보다 상위 레벨로 이해해요. 따라서 dbt-fabricspark를 실행할 때는 노드 config나 타깃 프로필에서 database를 사용하거나 설정해서는 안 돼요. 게다가 어댑터는 Lakehouse 내 스키마도 지원하지 않아요.

기본 파일 형식 구성 (Default file format configurations)

snapshotsmerge incremental 전략 같은 고급 incremental 전략 기능에 접근하려면, 모델을 table로 materialize할 때 Delta 파일 형식을 기본 파일 형식으로 사용하면 돼요.

프로젝트 파일에 최상위 구성을 설정하면 아주 편리해요:

dbt_project.yml

models:
  +file_format: delta
  
seeds:
  +file_format: delta
  
snapshots:
  +file_format: delta

각주 (Footnotes)

  1. location_root를 구성하면 dbt가 create table 문에서 위치 경로를 지정해요. 그러면 Fabric Lakehouse에서 테이블이 "managed"에서 "external"로 바뀌어요.

더 알아보기 (Learn more)