Confluent Cloud 설정

Confluent Cloud 설정

dbt-confluent 어댑터가 Confluent Cloud for Apache Flink에서 모델을 만들 때 지원하는 materialization과 설정들을 다루는 페이지예요. 이 어댑터는 배치 지향적인 기존 dbt와 달리, 지속적으로 실행되는 스트리밍 리소스를 만든다는 점이 특징이에요.

출처: 문서

본문

Materializations

Materialization 설명
view Flink SQL view를 drop하고 다시 생성해요.
streaming_table Flink SQL table과, 쿼리 결과를 테이블에 계속 쓰는 별도의 장기 실행 INSERT INTO 문을 만들어요. 테이블이 이미 있으면 어댑터가 스키마 drift를 확인하고 생성을 건너뛰어요. drop 후 재생성하려면 --full-refresh를 사용하세요.
streaming_source 커넥터(예: 목 데이터 생성을 위한 faker)로 뒷받침되는 Flink SQL table을 만들어요. connector config가 필요해요. 모델 SQL이 컬럼 정의를 정의해요. 테이블이 이미 있으면 어댑터가 스키마 drift를 확인하고 생성을 건너뛰어요. drop 후 재생성하려면 --full-refresh를 사용하세요.
ephemeral Flink에 materialize되지 않는 표준 dbt CTE 기반 쿼리 조각이에요.

지원되지 않는 materialization

  • table: 공식적으로 지원되지 않아요. 곧 지원될 예정이에요.
  • materialized_view: 지원되지 않아요. streaming_table을 사용하세요.
  • incremental: 지원되지 않아요. dbt의 배치 incremental 의미론은 Flink의 연속 처리 모델에 매핑되지 않아요. streaming_table을 사용하세요.
  • snapshot: 지원되지 않아요. Flink SQL에는 dbt snapshot이 필요한 배치 연산(MERGE, UPDATE)이 없어요.

Materialization별 구성

streaming_table

streaming_table materialization은 Kafka 토픽으로 뒷받침되는 테이블과 계속 실행되는 INSERT INTO 문을 만들어요.

models/my_streaming_model.sql

{{
  config(
    materialized='streaming_table',
    with={
      'changelog.mode': 'upsert',
      'kafka.retention.time': '7 d'
    }
  )
}}

SELECT
  order_id,
  customer_id,
  total_amount
FROM {{ ref('raw_orders') }}
WHERE status = 'completed'

with config

Config 타입 설명
with dict CREATE TABLE 문의 WITH 절에 전달되는 Confluent Cloud Flink SQL 테이블 옵션의 딕셔너리예요. 일반적인 옵션으로는 changelog.mode(append, upsert, retract), kafka.retention.time, key.format, value.format, scan.startup.mode 등이 있어요. 전체 목록은 CREATE TABLE WITH options 참조를 확인하세요.

streaming_source

streaming_source materialization은 커넥터로 뒷받침되는 테이블을 만들어요. connector config는 필수예요. 모델 SQL은 (SELECT 쿼리 대신) 컬럼 정의를 정의해요. Confluent Cloud에서 유효한 커넥터 값에는 faker(목 데이터 생성)와 AI 검색용 외부 테이블 커넥터가 있어요. 사용 가능한 커넥터와 옵션은 Confluent connector catalogFlink CREATE TABLE documentation을 참고하세요.

models/my_fake_orders.sql

{{
  config(
    materialized='streaming_source',
    connector='faker',
    with={
      'rows-per-second': '1',
      'number-of-rows': '100',
      'changelog.mode': 'append',
    }
  )
}}

order_id BIGINT,
price DECIMAL(10, 2),
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND,
PRIMARY KEY(order_id) NOT ENFORCED

connector config

Config 타입 필수 설명
connector string Yes 소스 테이블의 커넥터 유형이에요. Confluent Cloud에서 유효한 값에는 faker(목 데이터)와 외부 AI 검색 커넥터가 있어요.
with dict No 커넥터와 함께 WITH 절에 전달되는 추가 테이블 옵션이에요. 유효한 옵션은 커넥터 유형에 따라 달라져요.

상태 유지 동작과 --full-refresh

Confluent Cloud Flink SQL 테이블은 상태를 유지하는 장기 실행 리소스예요. streaming_tablestreaming_source materialization은 전통적인 배치 지향 dbt materialization과 다르게 동작해요:

  • 첫 번째 실행: 테이블이 생성되고(streaming_table의 경우) 계속 실행되는 INSERT INTO 문이 테이블을 채우기 시작해요.
  • --full-refresh 없는 후속 실행: 테이블이 이미 있으면 어댑터가 기존 컬럼 이름, 데이터 타입, WITH 옵션을 모델과 비교해요. drift가 없으면, 상태가 쌓였거나 다운스트림 소비자가 있는 테이블을 drop하는 것을 피하기 위해 실행이 모델을 건너뛰어요. drift가 감지되면 컴파일 오류로 실행이 실패해요. drift 감지는 모델별로 config(on_schema_drift='ignore')로 비활성화할 수 있어요.
  • --full-refresh 포함 실행: 기존 테이블을 drop하고 처음부터 다시 만들어 모든 데이터를 재처리해요.

테이블 스키마를 바꾸거나, WITH 옵션을 수정하거나, 데이터를 처음부터 다시 처리해야 할 땐 --full-refresh를 사용하세요:

dbt run --full-refresh --select my_streaming_model

알려진 제한 사항

  • 스키마 관리 없음: 어댑터는 스키마(Kafka 클러스터)나 데이터베이스(환경)를 만들거나 drop할 수 없어요. 이것들은 Confluent Cloud에서 관리해야 해요.
  • 테이블 이름 변경 없음: Flink SQL에서는 ALTER TABLE RENAME이 지원되지 않아요.
  • 비트랜잭션: Confluent Cloud Flink SQL은 트랜잭션을 지원하지 않아요. BEGINCOMMIT은 no-op이에요.
  • Seeds는 full_refresh 필요: 어댑터는 기본적으로 seeds에 대해 full_refresh: true를 설정해요.

더 알아보기 (Learn more)