DeltaStream 리소스 구성

DeltaStream 리소스 구성 (DeltaStream resource configurations)

이 문서는 dbt-DeltaStream 어댑터가 지원하는 다양한 리소스 구성을 설명해요. DeltaStream 고유의 구체화(materialization) 유형인 stream과 changelog, 그리고 인프라 구성 요소(store, entity, compute_pool, function 등)를 관리하는 방법을 강사 목소리로 살펴볼게요. SQL 모델 구성과 YAML 전용 리소스 구성의 차이, 시드(seed), 쿼리 관리 매크로도 함께 익힐 수 있습니다.

출처: 문서

본문

지원되는 구체화(Materializations)

DeltaStream은 스트리밍 처리 기능에 맞춘 몇 가지 고유한 구체화 유형을 지원합니다.

표준 구체화

Materialization 설명
ephemeral 이 구체화는 내부적으로 DeltaStream의 공통 테이블 표현식(common table expressions)을 사용합니다.
table 기존의 배치(batch) 테이블 구체화
materialized_view 기본 데이터가 변경될 때 자동으로 새로고침되는 지속 업데이트 뷰

스트리밍 구체화

Materialization 설명
stream 데이터를 실시간으로 처리하는 순수 스트리밍 변환
changelog 데이터의 변경을 추적하는 CDC(Change data capture) 스트림

인프라 구체화

Materialization 설명
store 외부 시스템 연결(Kafka, PostgreSQL 등)
entity store 내의 엔티티 정의
database 데이터베이스 정의
compute_pool 리소스 관리를 위한 컴퓨트 풀 정의
function Java의 사용자 정의 함수(UDF)
function_source UDF용 JAR 파일 소스
descriptor_source Protocol buffer 스키마 소스
schema_registry 스키마 레지스트리 연결(Confluent 등)

SQL 모델 구성

테이블 구체화

집계된 데이터를 위한 전통적인 배치 테이블을 만듭니다.

프로젝트 YAML 파일 구성:

models:
  <resource-path>:
    +materialized: table

SQL 구성:

{{ config(materialized = "table") }}

SELECT 
    date,
    SUM(amount) as daily_total
FROM {{ ref('transactions') }}
GROUP BY date

스트림 구체화

연속 스트리밍 변환을 만듭니다.

프로젝트 YAML 파일 구성:

models:
  <resource-path>:
    +materialized: stream
    +parameters:
      topic: 'stream_topic'
      value.format: 'json'
      key.format: 'primitive'
      key.type: 'VARCHAR'
      timestamp: 'event_time'

SQL 구성:

{{ config(
    materialized='stream',
    parameters={
        'topic': 'purchase_events',
        'value.format': 'json',
        'key.format': 'primitive',
        'key.type': 'VARCHAR',
        'timestamp': 'event_time'
    }
) }}

SELECT 
    event_time,
    user_id,
    action
FROM {{ ref('source_stream') }}
WHERE action = 'purchase'
스트림 구성 옵션
옵션 설명 필수?
materialized 모델이 어떻게 구체화될지. 스트리밍 모델을 만들려면 stream이어야 합니다. 필수
topic 스트림 출력의 토픽 이름. 필수
value.format 스트림 값의 형식('json', 'avro' 등). 필수
key.format 스트림 키의 형식('primitive', 'json' 등). 선택
key.type 스트림 키의 데이터 타입('VARCHAR', 'BIGINT' 등). 선택
timestamp 이벤트 타임스탬프로 사용할 컬럼 이름. 선택

체인지로그(Changelog) 구체화

데이터 스트림의 변경을 캡처합니다.

프로젝트 YAML 파일 구성:

models:
  <resource-path>:
    +materialized: changelog
    +parameters:
      topic: 'changelog_topic'
      value.format: 'json'
    +primary_key: [column_name]

SQL 구성:

{{ config(
    materialized='changelog',
    parameters={
        'topic': 'order_updates',
        'value.format': 'json'
    },
    primary_key=['order_id']
) }}

SELECT 
    order_id,
    status,
    updated_at
FROM {{ ref('orders_stream') }}
체인지로그 구성 옵션
옵션 설명 필수?
materialized 모델이 어떻게 구체화될지. 체인지로그 모델을 만들려면 changelog여야 합니다. 필수
topic 체인지로그 출력의 토픽 이름. 필수
value.format 체인지로그 값의 형식('json', 'avro' 등). 필수
primary_key 변경 추적을 위해 행을 고유하게 식별하는 컬럼 이름 목록. 필수

구체화된 뷰

지속적으로 업데이트되는 뷰를 만듭니다.

SQL 구성:

{{ config(materialized='materialized_view') }}

SELECT 
    product_id,
    COUNT(*) as purchase_count
FROM {{ ref('purchase_events') }}
GROUP BY product_id

YAML 전용 리소스 구성

DeltaStream은 인프라 구성 요소에 대해 두 가지 유형의 모델 정의를 지원합니다.

  1. 관리 리소스(Managed Resources, 모델) - dbt DAG에 자동으로 포함됩니다.
  2. 비관리 리소스(Unmanaged Resources, 소스) - 특정 매크로를 사용해 온디맨드로 생성됩니다.

관리 리소스와 비관리 리소스 중 무엇을 사용해야 하나?

  • 모든 인프라를 다른 환경에서 재현하거나, 그래프 연산자를 사용해 특정 리소스와 다운스트림 변환의 생성만 실행하려면 관리 리소스를 사용하세요.
  • 그렇지 않으면 플레이스홀더 파일을 피하기 위해 비관리 리소스를 사용하는 것이 더 간단할 수 있어요.

관리 리소스(모델)

관리 리소스는 dbt DAG에 자동으로 포함되며 모델로 정의됩니다.

version: 2
models:
  - name: my_kafka_store
    config:
      materialized: store
      parameters:
        type: KAFKA
        access_region: "AWS us-east-1"
        uris: "kafka.broker1.url:9092,kafka.broker2.url:9092"
        tls.ca_cert_file: "@/certs/us-east-1/self-signed-kafka-ca.crt"
  
  - name: ps_store
    config:
      materialized: store
      parameters:
        type: POSTGRESQL
        access_region: "AWS us-east-1"
        uris: "postgresql://mystore.com:5432/demo"
        postgres.username: "user"
        postgres.password: "password"
  
  - name: user_events_stream
    config:
      materialized: stream
      columns:
        event_time:
          type: TIMESTAMP
          not_null: true
        user_id:
          type: VARCHAR
        action:
          type: VARCHAR
      parameters:
        topic: 'user_events'
        value.format: 'json'
        key.format: 'primitive'
        key.type: 'VARCHAR'
        timestamp: 'event_time'
  
  - name: order_changes
    config:
      materialized: changelog
      columns:
        order_id:
          type: VARCHAR
          not_null: true
        status:
          type: VARCHAR
        updated_at:
          type: TIMESTAMP
      primary_key:
        - order_id
      parameters:
        topic: 'order_updates'
        value.format: 'json'
  
  - name: pv_kinesis
    config:
      materialized: entity
      store: kinesis_store
      parameters:
        'kinesis.shards': 3
  
  - name: my_compute_pool
    config:
      materialized: compute_pool
      parameters:
        'compute_pool.size': 'small'
        'compute_pool.timeout_min': 5
  
  - name: my_function_source
    config:
      materialized: function_source
      parameters:
        file: '@/path/to/my-functions.jar'
        description: 'Custom utility functions'
  
  - name: my_descriptor_source
    config:
      materialized: descriptor_source
      parameters:
        file: '@/path/to/schemas.desc'
        description: 'Protocol buffer schemas for data structures'
  
  - name: my_custom_function
    config:
      materialized: function
      parameters:
        args:
          - name: input_text
            type: VARCHAR
        returns: VARCHAR
        language: JAVA
        source.name: 'my_function_source'
        class.name: 'com.example.TextProcessor'
  
  - name: my_schema_registry
    config:
      materialized: schema_registry
      parameters:
        type: "CONFLUENT"
        access_region: "AWS us-east-1"
        uris: "https://url.to.schema.registry.listener:8081"
        'confluent.username': 'fake_username'
        'confluent.password': 'fake_password'
        'tls.client.cert_file': '@/path/to/tls/client_cert_file'
        'tls.client.key_file': '@/path/to/tls_key'

참고: 현재 dbt 제한 때문에 관리 YAML 전용 리소스에는 SELECT 문을 포함하지 않는 플레이스홀더 .sql 파일이 필요합니다. 예를 들어 my_kafka_store.sql을 다음 내용으로 만들면 됩니다.

-- Placeholder

비관리 리소스(소스)

비관리 리소스는 소스로 정의되며 특정 매크로를 사용해 온디맨드로 생성됩니다.

version: 2
sources:
  - name: infrastructure
    tables:
      - name: my_kafka_store
        config:
          materialized: store
          parameters:
            type: KAFKA
            access_region: "AWS us-east-1"
            uris: "kafka.broker1.url:9092,kafka.broker2.url:9092"
            tls.ca_cert_file: "@/certs/us-east-1/self-signed-kafka-ca.crt"
      
      - name: ps_store
        config:
          materialized: store
          parameters:
            type: POSTGRESQL
            access_region: "AWS us-east-1"
            uris: "postgresql://mystore.com:5432/demo"
            postgres.username: "user"
            postgres.password: "password"
      
      - name: user_events_stream
        config:
          materialized: stream
          columns:
            event_time:
              type: TIMESTAMP
              not_null: true
            user_id:
              type: VARCHAR
            action:
              type: VARCHAR
          parameters:
            topic: 'user_events'
            value.format: 'json'
            key.format: 'primitive'
            key.type: 'VARCHAR'
            timestamp: 'event_time'
      
      - name: order_changes
        config:
          materialized: changelog
          columns:
            order_id:
              type: VARCHAR
              not_null: true
            status:
              type: VARCHAR
            updated_at:
              type: TIMESTAMP
          primary_key:
            - order_id
          parameters:
            topic: 'order_updates'
            value.format: 'json'
      
      - name: pv_kinesis
        config:
          materialized: entity
          store: kinesis_store
          parameters:
            'kinesis.shards': 3
      
      - name: compute_pool_small
        config:
          materialized: compute_pool
          parameters:
            'compute_pool.size': 'small'
            'compute_pool.timeout_min': 5
      
      - name: my_function_source
        config:
          materialized: function_source
          parameters:
            file: '@/path/to/my-functions.jar'
            description: 'Custom utility functions'
      
      - name: my_descriptor_source
        config:
          materialized: descriptor_source
          parameters:
            file: '@/path/to/schemas.desc'
            description: 'Protocol buffer schemas for data structures'
      
      - name: my_custom_function
        config:
          materialized: function
          parameters:
            args:
              - name: input_text
                type: VARCHAR
            returns: VARCHAR
            language: JAVA
            source.name: 'my_function_source'
            class.name: 'com.example.TextProcessor'
      
      - name: my_schema_registry
        config:
          materialized: schema_registry
          parameters:
            type: "CONFLUENT"
            access_region: "AWS us-east-1"
            uris: "https://url.to.schema.registry.listener:8081"
            'confluent.username': 'fake_username'
            'confluent.password': 'fake_password'
            'tls.client.cert_file': '@/path/to/tls/client_cert_file'
            'tls.client.key_file': '@/path/to/tls_key'

비관리 리소스를 만들려면:

# 모든 소스 생성
dbt run-operation create_sources

# 특정 소스 생성
dbt run-operation create_source_by_name --args '{source_name: infrastructure}'

Store 구성

Kafka store

- name: my_kafka_store
  config:
    materialized: store
    parameters:
      type: KAFKA
      access_region: "AWS us-east-1"
      uris: "kafka.broker1.url:9092,kafka.broker2.url:9092"
      tls.ca_cert_file: "@/certs/us-east-1/self-signed-kafka-ca.crt"

PostgreSQL store

- name: postgres_store
  config:
    materialized: store
    parameters:
      type: POSTGRESQL
      access_region: "AWS us-east-1"
      uris: "postgresql://mystore.com:5432/demo"
      postgres.username: "user"
      postgres.password: "password"

Entity 구성

- name: kinesis_entity
  config:
    materialized: entity
    store: kinesis_store
    parameters:
      'kinesis.shards': 3

Compute pool 구성

- name: processing_pool
  config:
    materialized: compute_pool
    parameters:
      'compute_pool.size': 'small'
      'compute_pool.timeout_min': 5

리소스 참조하기

관리 리소스

표준 ref() 함수를 사용하세요.

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

비관리 리소스

source() 함수를 사용하세요.

SELECT * FROM {{ source('infrastructure', 'user_events_stream') }}

시드(Seeds)

seed 구체화를 사용해 CSV 데이터를 기존 DeltaStream 엔티티에 로드합니다. 새 테이블을 만드는 기존 dbt 시드와 달리 DeltaStream 시드는 데이터를 이미 존재하는 엔티티에 삽입합니다.

구성

시드는 YAML에서 다음 속성으로 구성해야 합니다.

필수:

  • entity: 데이터를 삽입할 대상 엔티티의 이름

선택:

  • store: 엔티티가 있는 store의 이름(엔티티가 store에 없으면 생략)
  • with_params: WITH 절용 파라미터 사전
  • quote_columns: 따옴표로 묶을 컬럼 제어. 기본값: false(따옴표 없음). 값:
    • true: 모든 컬럼을 따옴표로 묶음
    • false: 컬럼을 따옴표로 묶지 않음(기본값)
    • string: '*'로 설정하면 모든 컬럼을 따옴표로 묶음
    • list: 따옴표로 묶을 컬럼 이름 목록

예시 구성

Store 사용(따옴표 활성화):

# seeds.yml
version: 2

seeds:
  - name: user_data_with_store_quoted
    config:
      entity: 'user_events'
      store: 'kafka_store'
      with_params:
        kafka.topic.retention.ms: '86400000'
        partitioned: true
      quote_columns: true  # 모든 컬럼을 따옴표로 묶음

사용법

  1. seeds/ 디렉터리에 CSV 파일을 배치하세요.
  2. 필수 entity 파라미터로 YAML에서 시드를 구성하세요.
  3. 엔티티가 store에 있으면 선택적으로 store를 지정하세요.
  4. dbt seed를 실행해 데이터를 로드하세요.

중요: 시드를 실행하기 전에 대상 엔티티가 DeltaStream에 이미 존재해야 합니다. 시드는 데이터만 삽입하며 엔티티를 만들지는 않습니다.

함수 및 소스 구체화

DeltaStream은 특수 구체화를 통해 사용자 정의 함수(UDF)와 그 의존성을 지원합니다.

파일 첨부 지원

어댑터는 함수 소스와 descriptor 소스에 대한 원활한 파일 첨부를 제공합니다.

  • 표준화된 인터페이스: 함수 소스와 descriptor 소스 모두에 공통 파일 처리 로직
  • 경로 해석: 절대 경로와 상대 경로 모두 지원(프로젝트 상대 경로용 @ 구문 포함)
  • 자동 검증: 첨부 전에 파일이 존재하고 접근 가능한지 검증

함수 소스(Function source)

Java 함수를 담은 JAR 파일에서 함수 소스를 만듭니다.

SQL 구성:

{{ config(
    materialized='function_source',
    parameters={
        'file': '@/path/to/my-functions.jar',
        'description': 'Custom utility functions'
    }
) }}

SELECT 1 as placeholder

Descriptor 소스

컴파일된 protocol buffer descriptor 파일에서 descriptor 소스를 만듭니다.

SQL 구성:

{{ config(
    materialized='descriptor_source',
    parameters={
        'file': '@/path/to/schemas.desc',
        'description': 'Protocol buffer schemas for data structures'
    }
) }}

SELECT 1 as placeholder

참고: descriptor 소스는 .proto 원시 파일이 아닌 컴파일된 .desc 파일이 필요합니다. protobuf 스키마를 protoc --descriptor_set_out=schemas/my_schemas.desc schemas/my_schemas.proto로 컴파일하세요.

함수(Function)

함수 소스를 참조하는 사용자 정의 함수를 만듭니다.

SQL 구성:

{{ config(
    materialized='function',
    parameters={
        'args': [
            {'name': 'input_text', 'type': 'VARCHAR'}
        ],
        'returns': 'VARCHAR',
        'language': 'JAVA',
        'source.name': 'my_function_source',
        'class.name': 'com.example.TextProcessor'
    }
) }}

SELECT 1 as placeholder

스키마 레지스트리(Schema registry)

스키마 레지스트리 연결을 만듭니다.

SQL 구성:

{{ config(
    materialized='schema_registry',
    parameters={
        'type': 'CONFLUENT',
        'access_region': 'AWS us-east-1',
        'uris': 'https://url.to.schema.registry.listener:8081',
        'confluent.username': 'fake_username',
        'confluent.password': 'fake_password',
        'tls.client.cert_file': '@/path/to/tls/client_cert_file',
        'tls.client.key_file': '@/path/to/tls_key'
    }
) }}

SELECT 1 as placeholder

쿼리 관리 매크로

DeltaStream dbt 어댑터는 dbt에서 직접 실행 중인 쿼리를 관리하고 종료하는 매크로를 제공합니다.

모든 쿼리 나열

list_all_queries 매크로는 DeltaStream이 현재 알고 있는 모든 쿼리를 상태, 소유자, SQL과 함께 표시합니다.

dbt run-operation list_all_queries

쿼리 설명

describe_query 매크로를 사용해 특정 쿼리의 로그와 세부 정보를 확인하세요.

dbt run-operation describe_query --args '{query_id: "<QUERY_ID>"}'

특정 쿼리 종료

terminate_query 매크로로 ID를 사용해 쿼리를 종료합니다.

dbt run-operation terminate_query --args '{query_id: "<QUERY_ID>"}'

모든 실행 중인 쿼리 종료

terminate_all_queries 매크로로 현재 실행 중인 모든 쿼리를 종료합니다.

dbt run-operation terminate_all_queries

쿼리 재시작

restart_query 매크로로 ID를 사용해 실패한 쿼리를 재시작합니다.

dbt run-operation restart_query --args '{query_id: "<QUERY_ID>"}'

애플리케이션 매크로

여러 문을 단일 단위로 실행

application 매크로는 여러 DeltaStream SQL 문을 all-or-nothing 시맨틱의 단일 작업 단위로 실행할 수 있게 해줍니다.

dbt run-operation application --args '{
  application_name: "my_data_pipeline",
  statements: [
    "USE DATABASE my_db",
    "CREATE STREAM user_events WITH (topic='"'events'"'', value.format='"'json'"')",
    "CREATE MATERIALIZED VIEW user_counts AS SELECT user_id, COUNT(*) FROM user_events GROUP BY user_id"
  ]
}'

문제 해결

함수 소스 준비 상태

함수를 만들 때 "function source is not ready" 오류가 발생하면:

  1. 자동 재시도: 어댑터가 지수 백오프로 함수 생성을 자동 재시도합니다.
  2. 타임아웃 구성: 기본 30초 타임아웃은 큰 JAR 파일에 필요하면 연장할 수 있어요.
  3. 의존성 순서: 함수 소스가 의존 함수보다 먼저 생성되었는지 확인하세요.
  4. 수동 재시도: 자동 재시도가 실패하면 몇 분 기다린 후 작업을 재시도하세요.

파일 첨부 문제

함수 소스와 descriptor 소스에서 파일 첨부 문제가 있으면:

  1. 파일 경로: 프로젝트 상대 경로에는 @/path/to/file 구문을 사용하세요.
  2. 파일 유형:
    • 함수 소스는 .jar 파일이 필요합니다.
    • Descriptor 소스는 컴파일된 .desc 파일이 필요합니다(.proto 아님).
  3. 파일 검증: 어댑터가 첨부를 시도하기 전에 파일 존재를 검증합니다.
  4. 컴파일: Descriptor 소스의 경우 protobuf 파일이 컴파일되었는지 확인하세요: protoc --descriptor_set_out=output.desc input.proto

더 알아보기 (Learn more)

DeltaStream 어댑터는 스트리밍 데이터 처리에 특화된 dbt 확장이에요. dbt의 일반적인 리소스 구성(resource configs)과 점진적 전략 개념을 먼저 익힌 뒤, 다른 클라우드 어댑터(BigQuery, Databricks, Snowflake)와 비교하며 학습하면 어댑터별 추상화가 훨씬 잘 이해됩니다.