소스(Source)로 DAG에 데이터 입력 정의하기

소스(Source)로 DAG에 데이터 입력 정의하기

dbt 프로젝트에서 가장 먼저 만나는 데이터는 보통 Extract & Load 도구가 웨어하우스에 쌓아 둔 로우 데이터예요. dbt는 이 테이블들을 **소스(source)**로 선언해서 이름을 붙이고, 모델과의 관계까지 정의할 수 있어요. 소스를 선언하면 모델에서 {{ source() }} 함수로 참조해 데이터 계보(lineage)를 만들고, 소스 데이터에 대한 가정을 테스트하며, 데이터의 신선도(freshness)까지 계산할 수 있어요.

출처: dbt 공식 문서 — Add sources to your DAG

소스 선언하기

소스는 .yml 파일의 sources: 키 아래에 정의해요. 예시를 볼게요.


sources:
  - name: jaffle_shop
    database: raw  
    schema: jaffle_shop  
    tables:
      - name: orders
      - name: customers

  - name: stripe
    tables:
      - name: payments

기본적으로 schemaname과 같아져요. 기존 스키마와 다른 이름을 쓰고 싶을 때만 schema를 추가하면 돼요. 이 파일들에 익숙하지 않다면 properties.yml 파일 문서를 먼저 확인해 두는 게 좋아요.

소스에서 데이터 선택하기

소스를 정의하고 나면 {{ source() }} 함수로 모델에서 참조할 수 있어요.

select
  ...

from {{ source('jaffle_shop', 'orders') }}

left join {{ source('jaffle_shop', 'customers') }} using (customer_id)

dbt는 이걸 전체 테이블 이름으로 컴파일해요.


select
  ...

from raw.jaffle_shop.orders

left join raw.jaffle_shop.customers using (customer_id)

{{ source() }} 함수를 쓰면 모델과 소스 테이블 사이에 의존 관계가 생겨요. 이 의존 관계가 DAG에 그대로 나타나서, 모델이 어떤 소스에 의존하는지 한눈에 보이게 돼요.

소스 테스트하고 문서화하기

소스에도 모델처럼 데이터 테스트를 추가하고, 설명(description)을 붙여 문서화 사이트에 렌더링되게 할 수 있어요.


sources:
  - name: jaffle_shop
    description: This is a replica of the Postgres database used by our app
    tables:
      - name: orders
        database: raw
        description: >
          One record per order. Includes cancelled and deleted orders.
        columns:
          - name: id
            description: Primary key of the orders table
            data_tests:
              - unique
              - not_null
          - name: status
            description: Note that the status can change over time

      - name: ...

  - name: ...

소스에 사용할 수 있는 속성의 전체 목록은 reference section에서 확인할 수 있어요.

소스 데이터 신선도(freshness)

설정을 몇 개 추가하면 dbt가 소스 테이블 데이터의 "신선도"를 선택적으로 캡처할 수 있어요. 데이터 파이프라인이 건강한 상태인지 이해하는 데 유용하고, 웨어하우스의 서비스 수준 계약(SLA)을 정의하는 핵심 요소가 돼요.

소스 신선도 선언하기

소스 신선도 정보를 설정하려면 소스에 freshness 블록을, 테이블 선언에는 loaded_at_field를 추가해요.


sources:
  - name: jaffle_shop
    database: raw
    config: 
      freshness: # default freshness
        # changed to config in v1.9
        warn_after: {count: 12, period: hour}
        error_after: {count: 24, period: hour}
      loaded_at_field: _etl_loaded_at # changed to config in v1.10

    tables:
      - name: orders
        config:
          freshness: # make this a little more strict
            warn_after: {count: 6, period: hour}
            error_after: {count: 12, period: hour}

      - name: customers # this inherits the default freshness defined in the jaffle_shop source block at the beginning


      - name: product_skus
        config:
          freshness: null # do not check freshness for this table

freshness 블록에서는 warn_aftererror_after 중 하나 또는 둘 다 제공할 수 있어요. 둘 다 없으면 dbt는 이 소스의 테이블에 대해 신선도를 계산하지 않아요. 추가로 테이블의 신선도를 계산하려면 loaded_at_field가 필수예요(웨어하우스 메타데이터로 신선도를 계산할 수 있는 경우는 제외). 이 설정은 계층적으로 적용돼서, 소스에 지정한 freshnessloaded_at_field 값이 그 소스에 정의된 모든 테이블로 흘러가요. 모든 테이블의 loaded_at_field가 같다면 최상위 소스 선언에 한 번만 지정하면 돼요.

소스 신선도 확인하기

신선도 정보를 얻으려면 dbt source freshness 명령을 실행해요.

$ dbt source freshness

내부적으로 dbt는 신선도 속성을 이용해 아래와 같은 select 쿼리를 만드는데, 이 쿼리는 쿼리 로그에서 찾을 수 있어요.

select
  max(_etl_loaded_at) as max_loaded_at,
  convert_timezone('UTC', current_timestamp()) as calculated_at
from raw.jaffle_shop.orders

이 쿼리의 결과로 소스가 신선한지 아닌지가 결정돼요.

소스 신선도 기준으로 모델 빌드하기

dbt를 쓰면서 가장 좋은 방법으로 권장하는 것은 소스 신선도 설정을 failure_after 같은 속성 대신 .yml 파일에 두고 모델 레벨에서 정의하는 거예요. 소스 신선도를 기준으로 모델을 빌드하려면 다음 순서를 따라요.

  1. dbt source freshness를 실행해 소스의 신선도를 확인해요.
  2. dbt build --select source_status:fresher+ 명령으로 더 신선한 소스 하위의 모델을 빌드하고 테스트해요.

이 명령을 순서대로 쓰면 모델이 항상 최신 데이터로 갱신돼요. 바뀌지 않은 데이터에 낭비되는 컴퓨팅을 없애고, 필요한 때에만 모델을 빌드해요. 소스 신선도 체크를 30분으로 설정하고, 1시간마다 다시 빌드하는 잡을 돌려 보세요. 소스 신선도가 만료되면 해당 모델들을 한 번에 다시 빌드하는 구조예요. 자세한 내용은 Source freshness check frequency 문서를 참고해요.

필터(filter)

테이블에 따라서는 전체 스캔을 막기 위해 특정 컬럼에 필터가 필요한 경우가 있어요. 비용이 클 수 있으니까요. 이런 테이블의 신선도 체크를 하려면 설정에 filter 인자를 추가하면 돼요. 예를 들어 filter: _etl_loaded_at >= date_sub(current_date(), interval 1 day)처럼요. 위 예시에서 결과 쿼리는 이렇게 돼요.

select
  max(_etl_loaded_at) as max_loaded_at,
  convert_timezone('UTC', current_timestamp()) as calculated_at
from raw.jaffle_shop.orders
where _etl_loaded_at >= date_sub(current_date(), interval 1 day)

더 알아보기