Kafka용 Openflow 커넥터 설정

Kafka용 Openflow 커넥터 설정

이 주제는 Openflow Connector for Kafka를 설정하는 단계를 설명해요.

출처: Snowflake 문서

본문

참고: 이 커넥터는 Snowflake Connector Terms가 적용돼요.

전제 조건(Prerequisites)

  1. Snowflake Openflow Connector for Kafka을 검토했는지 확인해요.
  2. Set up Openflow - BYOC 또는 Set up Openflow - Snowflake Deployments를 설정했는지 확인해요.
  3. Openflow - Snowflake Deployments를 사용한다면 configuring required domains을 검토하고, Kafka 커넥터의 필수 도메인에 대한 액세스를 부여했는지 확인해요. 커넥터는 클러스터의 모든 Kafka 브로커에 연결할 수 있어야 해요.

Snowflake 계정 설정하기

Snowflake 계정 관리자로서 다음 작업을 수행해요:

  1. 타입이 SERVICE인 새 Snowflake 서비스 사용자를 만들어요.

  2. 새 역할을 만들거나 기존 역할을 사용하고 데이터베이스 권한(Database privileges)을 부여해요. 커넥터는 대상 테이블을 만들려면 사용자가 필요해요. 사용자가 Snowflake 객체 관리를 위한 필요한 권한을 갖고 있는지 확인해요:

    객체 권한 비고
    Database USAGE
    Schema USAGE
    Table OWNERSHIP 커넥터가 테이블에 데이터를 수집하는 데 필요

    Snowflake는 더 나은 액세스 제어를 위해 각 Kafka 클러스터마다 별도의 사용자와 역할을 만들 것을 권장해요. 다음 스크립트를 사용해 커스텀 역할을 만들고 구성할 수 있어요(SECURITYADMIN 또는 그에 해당하는 역할 필요):

    USE ROLE securityadmin;
    CREATE ROLE openflow_kafka_connector_role_1;
    
    GRANT USAGE ON DATABASE kafka_db TO ROLE openflow_kafka_connector_role_1;
    GRANT USAGE ON SCHEMA kafka_schema TO ROLE openflow_kafka_connector_role_1;
    

    참고: 권한은 커넥터 역할에 직접 부여되어야 하며 상속될 수 없어요.

  3. 대상 테이블을 구성해요. Snowflake는 스키마 변경에 서버 측 스키마 진화(server-side schema evolution)를, DML 오류 로깅용 오류 테이블(error table)을 사용할 것을 강력히 권장해요. 다음 예시는 테이블을 만들고 적절한 OWNERSHIP 권한을 추가하는 방법을 보여줘요.

    USE ROLE openflow_kafka_connector_role_1;
    
    CREATE TABLE kafka_db.kafka_schema.<DESTINATION_TABLE_NAME> (
      kafkaMetadata variant
    )
    ENABLE_SCHEMA_EVOLUTION = TRUE
    ERROR_LOGGING = TRUE;
    
    USE ROLE securityadmin;
    GRANT OWNERSHIP ON TABLE existing_table1 TO ROLE openflow_kafka_connector_role_1;
    

    커넥터는 자동 스키마 감지와 진화를 지원해요. Snowflake의 테이블 구조는 커넥터가 로드한 새 데이터의 구조를 지원하도록 자동으로 정의되고 진화해요. 커넥터는 레코드 콘텐츠의 1차 키(첫 번째 수준 키)를 이름으로(대소문자 무시) 테이블 컬럼에 자동으로 매핑해요. 스키마 진화를 활성화하면 Snowflake가 수신 스트림에서 감지된 새 컬럼을 추가해 대상 테이블을 자동으로 확장하고, 새 데이터 패턴을 수용하기 위해 NOT NULL 제약 조건을 해제할 수 있어요. 자세한 내용은 Table schema evolution을 참고해요. ENABLE_SCHEMA_EVOLUTION이 활성화되지 않았다면, 테이블 정의를 확장해 스키마를 수동으로 만들어야 해요. 커넥터는 레코드 콘텐츠의 1차 키를 이름으로 테이블 컬럼에 매칭하려고 시도해요. JSON의 키가 테이블 컬럼과 일치하지 않으면 커넥터는 그 키를 무시해요.

  4. (선택) 시크릿 매니저를 구성해요. Snowflake는 이 단계를 강력히 권장해요. Openflow가 지원하는 시크릿 매니저(예: AWS, Azure, Hashicorp)를 구성하고 공개·개인 키를 시크릿 저장소에 보관해요.

    • 구성 후 시크릿 매니저에 인증할 방법을 결정해요. AWS에서는 Openflow와 연결된 EC2 인스턴스 역할을 사용해 다른 시크릿을 저장할 필요가 없도록 하는 것을 권장해요.
    • Openflow에서 오른쪽 상단의 햄버거 메뉴에서 이 시크릿 매니저와 연결된 Parameter Provider를 구성해요. Controller Settings > Parameter Provider로 이동해 파라미터 값을 가져와요.
    • 민감한 값이 Openflow 안에 저장되지 않도록 모든 자격 증명을 연결된 파라미터 경로로 참조해요.

    사용자에게 액세스 부여 커넥터가 수집한 원시 데이터에 액세스해야 하는 다른 Snowflake 사용자(예: Snowflake에서 커스텀 처리용)에게 1단계에서 만든 역할을 부여해요.

커넥터 설정하기

데이터 엔지니어로서 커넥터를 설치하고 구성하려면 다음 작업을 수행해요:

커넥터 설치하기

커넥터를 설치하려면:

  1. Openflow의 Connector library 탭으로 이동해요.
  2. Openflow 커넥터 페이지에서 커넥터를 찾고 Install을 선택해요.
  3. Select runtime 대화상자에서 Available runtimes 드롭다운 목록에서 런타임을 선택하고 Add를 선택해요.

참고: 커넥터를 설치하기 전에, 수집된 데이터를 저장할 Snowflake에 데이터베이스, 스키마, 테이블을 만들었는지 확인해요.

  1. Snowflake 계정 자격 증명으로 배포에 인증하고, 런타임 애플리케이션이 Snowflake 계정에 접근하는 것을 허용하라는 프롬프트가 나오면 Allow를 선택해요. 커넥터 설치 과정은 완료되는 데 몇 분이 걸려요.
  2. Snowflake 계정 자격 증명으로 런타임에 인증해요. Openflow 캔버스에 커넥터 프로세스 그룹이 추가된 상태로 나타나요.

커넥터 구성하기

  1. 필요하면 내장 파라미터를 구성하기 전에 커넥터 구성을 커스터마이즈해요. 일부 일반적인 커스터마이즈에는 전용 가이드가 있어요(예: custom transformations, Avro 및 Protobuf 데이터 타입 수집, dead-letter-queue 처리). Snowflake CoCo의 Openflow 스킬로 커스터마이즈를 적용할 수도 있어요. 자세한 내용은 Configuring custom transformations을 참고해요.
  2. 프로세스 그룹 파라미터를 채워요.
    • 가져온 프로세스 그룹을 마우스 오른쪽 버튼으로 클릭하고 Parameters를 선택해요.
    • 필수 파라미터 값을 채워요.

파라미터

다음 표는 Kafka용 Openflow 커넥터의 파라미터를 설명해요:

파라미터 설명 필수
Kafka Auto Offset Reset 이전 컨슈머 오프셋이 발견되지 않을 때 적용되는 자동 오프셋 구성(Kafka auto.offset.reset 속성에 해당). 가능한 값: earliest(오프셋을 가장 이른 오프셋으로 자동 재설정), latest(오프셋을 가장 최근 오프셋으로 자동 재설정), none(소비자 그룹의 이전 오프셋이 없으면 컨슈머에 예외 발생). 기본값: latest 예
Kafka Bootstrap Servers 포트를 포함해야 하는, 쉼표로 구분된 Kafka 부트스트랩 서버 목록(예: kafka-broker:9092) 예
Kafka Consumer Group ID 커넥터가 사용하는 소비자 그룹의 ID. 임의일 수 있지만 고유해야 함 예
Kafka SASL Password SASL512 SCRAM 메커니즘 사용 시 구성된 비밀번호와 함께 제공되는 비밀번호
Kafka SASL Username SASL512 SCRAM 메커니즘 사용 시 구성된 비밀번호와 함께 제공되는 사용자 이름
Kafka Topic Format names / pattern 중 하나. 제공된 "Kafka Topics"가 쉼표로 구분된 이름 목록인지 단일 정규식인지 지정 예
Kafka Topics 쉼표로 구분된 Kafka 토픽 목록 또는 정규식 예
Snowflake Destination Database 데이터가 저장되는 데이터베이스. Snowflake에 이미 존재해야 함. 이름은 대소문자를 구분함. 인용되지 않은 식별자는 대문자로 제공 예
Snowflake Destination Schema 데이터가 저장되는 스키마. Snowflake에 이미 존재해야 함. 이름은 대소문자를 구분함. 인용되지 않은 식별자는 대문자로 제공. 다음 예시 참고: CREATE SCHEMA SCHEMA_NAME 또는 CREATE SCHEMA schema_name: SCHEMA_NAME 사용. CREATE SCHEMA "schema_name" 또는 CREATE SCHEMA "SCHEMA_NAME": 각각 schema_name 또는 SCHEMA_NAME 사용 예
Snowflake Destination Table 데이터가 저장되는 테이블. Snowflake에 이미 존재해야 함. 이름은 대소문자를 구분함. 인용되지 않은 식별자는 대문자로 제공 예

커넥터 시작하기

  1. 평면(plane)을 마우스 오른쪽 버튼으로 클릭하고 Enable all Controller Services를 선택해요.
  2. 평면을 마우스 오른쪽 버튼으로 클릭하고 Start를 선택해요. 커넥터가 데이터 수집을 시작해요.

KAFKAMETADATA 컬럼 이해하기

커넥터는 Kafka 레코드에 대한 메타데이터로 KAFKAMETADATA 구조를 채워요. 이 구조는 다음 정보를 포함해요:

필드 데이터 타입 설명
topic String 레코드가 온 Kafka 토픽의 이름
partition number 토픽 안의 파티션 번호. (이것은 Kafka 파티션이지 Snowflake 마이크로 파티션이 아님)
offset number 그 파티션 안의 오프셋
timestamp number 레코드가 Kafka에 추가된 시각
key String 메시지가 Kafka KeyedMessage이면 그 메시지의 키. 커넥터가 키를 RECORD_METADATA에 저장하려면 Kafka 구성 속성의 key.converter 파라미터가 org.apache.kafka.connect.storage.StringConverter로 설정되어야 함. 그렇지 않으면 커넥터가 키를 무시함
headers Object 헤더는 레코드와 연결된 사용자 정의 키-값 쌍. 각 레코드는 0, 1 또는 여러 개의 헤더를 가질 수 있음

수집 지연 시간 측정하기

변경 추적, 증분 처리, 행 수정 시간 기반의 시간 여행 쿼리에는 ROW_TIMESTAMP 기능을 사용할 수 있어요. 대상 테이블에서 다음 명령을 실행해 활성화해요:

ALTER TABLE <DESTINATION_TABLE> SET ROW_TIMESTAMP = TRUE;

행 타임스탬프를 활성화하면 테이블이 METADATA$ROW_LAST_COMMIT_TIME 컬럼을 노출하는데, 이 컬럼은 각 행이 마지막으로 수정된 시각을 반환해요. 자세한 내용은 METADATA$ROW_LAST_COMMIT_TIME을 참고해요.

참고: 행 타임스탬프는 interactive tables에서는 사용할 수 없어요. 자세한 내용은 Snowflake interactive analytics을 참고해요.

Apache Iceberg™ 테이블과 함께 커넥터 사용하기

커넥터는 Snowflake 관리형 Apache Iceberg™ 테이블에 데이터를 수집할 수 있어요. 커넥터는 Iceberg 테이블을 자동으로 만들지 않아요. 커넥터를 실행하기 전에 Iceberg 테이블을 수동으로 만들어야 해요.

커넥터는 표준 Snowflake 테이블과 같은 방식으로 Iceberg 대상 테이블의 서버 측 스키마 진화를 지원해요. 대상 테이블에 ENABLE_SCHEMA_EVOLUTION = TRUE가 있으면 Snowflake가 수신 스트림에서 감지된 새 컬럼을 자동으로 추가하고, 새 데이터 패턴을 수용하기 위해 NOT NULL 제약 조건을 해제해요. 스키마 진화 동작에 대한 자세한 내용은 Table schema evolution을 참고해요.

Iceberg 테이블은 다음 두 저장 옵션 중 하나를 사용할 수 있어요:

  • Snowflake storage: Snowflake가 Iceberg 테이블 파일을 저장하고 관리하므로 외부 볼륨을 만들거나 커넥터에 액세스를 부여할 필요가 없어요.
  • 외부 볼륨을 통해 접근하는, 여러분이 관리하는 외부 클라우드 저장소. 커넥터 역할에 외부 볼륨에 대한 USAGE를 부여해야 해요.

외부 볼륨에 USAGE 부여하기

이 단계는 Iceberg 테이블이 여러분이 관리하는 외부 볼륨을 사용할 때만 적용돼요. 테이블이 Snowflake storage를 사용하면 이 단계를 건너뛰어요. 예를 들어 Iceberg 테이블이 kafka_external_volume 외부 볼륨을 사용하고 커넥터가 openflow_kafka_connector_role 역할을 사용한다면 다음 문을 실행해요:

USE ROLE ACCOUNTADMIN;
GRANT USAGE ON EXTERNAL VOLUME kafka_external_volume TO ROLE openflow_kafka_connector_role;

수집용 Apache Iceberg™ 테이블 만들기

Iceberg 테이블을 만들 때 Iceberg 데이터 타입(VARIANT 포함) 또는 호환되는 Snowflake 타입을 사용할 수 있어요. 예를 들어 다음 메시지를 고려해 보죠:

{
  "id": 1,
  "name": "Steve",
  "body_temperature": 36.6,
  "approved_coffee_types": ["Espresso", "Doppio", "Ristretto", "Lungo"],
  "animals_possessed": {
    "dogs": true,
    "cats": false
  },
  "options": {
    "can_walk": true,
    "can_talk": false
  },
  "date_added": "2024-10-15"
}

예시 메시지에 대한 Iceberg 테이블을 만들려면 다음 문 중 하나를 사용해요. Snowflake storage를 사용하려면 EXTERNAL_VOLUME = 'SNOWFLAKE_MANAGED'로 설정하고 BASE_LOCATION은 생략해요:

CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
  kafkaMetadata OBJECT(
    topic STRING,
    partition INTEGER,
    offset BIGINT,
    key STRING,
    headers MAP(STRING, STRING),
    timestamp BIGINT
  ),
  id INT,
  name string,
  body_temperature float,
  approved_coffee_types array(string),
  animals_possessed variant,
  date_added date,
  options object(can_walk boolean, can_talk boolean)
)
EXTERNAL_VOLUME = 'SNOWFLAKE_MANAGED'
CATALOG = 'SNOWFLAKE'
ICEBERG_VERSION = 3;

자체 외부 볼륨을 사용하려면 EXTERNAL_VOLUME을 볼륨 이름으로 설정하고 BASE_LOCATION을 제공해요:

CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
  kafkaMetadata OBJECT(
    topic STRING,
    partition INTEGER,
    offset BIGINT,
    key STRING,
    headers MAP(STRING, STRING),
    timestamp BIGINT
  ),
  id INT,
  name string,
  body_temperature float,
  approved_coffee_types array(string),
  animals_possessed variant,
  date_added date,
  options object(can_walk boolean, can_talk boolean)
)
EXTERNAL_VOLUME = 'my_volume'
CATALOG = 'SNOWFLAKE'
BASE_LOCATION = 'my_location/my_iceberg_table'
ICEBERG_VERSION = 3;

Interactive Tables와 함께 커넥터 사용하기

Interactive tables는 저지연·고동시성 쿼리에 최적화된 특별한 타입의 Snowflake 테이블이에요. 자세한 내용은 Snowflake interactive analytics을 참고해요.

  1. interactive table을 만들어요:
    CREATE INTERACTIVE TABLE REALTIME_METRICS (
      metric_name VARCHAR,
      metric_value NUMBER,
      source_topic VARCHAR,
      timestamp TIMESTAMP_NTZ
    ) AS (SELECT
       $1:M_NAME::VARCHAR,
       $1:M_VALUE::NUMBER,
       $1:RECORD_METADATA.topic::VARCHAR,
       $1:RECORD_METADATA.timestamp::TIMESTAMP_NTZ
       from TABLE(DATA_SOURCE(TYPE => 'STREAMING')));
    

중요한 고려 사항:

  • Interactive tables에는 특정 제한과 쿼리 제약이 있어요. 커넥터와 함께 사용하기 전에 Snowflake interactive analytics을 검토해요.
  • Interactive tables의 경우 필요한 모든 변환이 테이블 정의에서 처리되어야 해요.
  • Interactive warehouses는 interactive tables를 효율적으로 쿼리하는 데 필요해요.

대상 테이블의 고객 정의 스키마와 함께 커넥터 사용하기

커넥터는 각 Kafka 레코드를 Snowflake 테이블에 삽입할 행으로 취급해요. 예를 들어 다음 JSON처럼 구조화된 메시지 내용이 있는 Kafka 토픽이 있다고 가정해 보죠:

{
  "order_id": 12345,
  "customer_name": "John",
  "order_total": 100.00,
  "isPaid": true
}

기본적으로 ENABLE_SCHEMA_EVOLUTION = TRUE 기능 덕분에 JSON의 모든 필드를 지정할 필요는 없어요. 하지만 정적 스키마를 선호한다면 다음을 실행해 만들 수 있어요:

CREATE TABLE ORDERS (
  kafkaMetadata OBJECT,
  order_id NUMBER,
  customer_name VARCHAR,
  order_total FLOAT,
  ispaid BOOLEAN
);

고객 정의 PIPE와 함께 커넥터 사용하기

자체 파이프를 만들기로 선택하면 파이프의 COPY INTO 문에 데이터 변환 로직을 정의할 수 있어요. 필요에 따라 컬럼 이름을 바꾸고 데이터 타입을 캐스트할 수 있어요. 예를 들어:

CREATE TABLE ORDERS (
  order_id VARCHAR,
  customer_name VARCHAR,
  order_total VARCHAR,
  ispaid VARCHAR
);

CREATE PIPE ORDERS AS
COPY INTO ORDERS
SELECT
  $1:order_id::STRING,
  $1:customer_name,
  $1:order_total::STRING,
  $1:isPaid::STRING
FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING'));

자체 파이프를 정의하면 대상 테이블 컬럼이 JSON 키와 일치할 필요가 없어요. 컬럼을 원하는 이름으로 바꾸고 필요하면 데이터 타입을 캐스트할 수 있어요.

커넥터가 커스텀 파이프로 작동하도록 조정하려면:

  1. Openflow 캔버스에서 Kafka 수집 플로우에 사용된 PublishSnowpipeStreaming 프로세서를 마우스 오른쪽 버튼으로 클릭해요.
  2. 컨텍스트 메뉴에서 Configure를 선택해요.
  3. Properties 탭으로 이동해요.
  4. Destination type 필드에서 Pipe를 선택해요.
  5. Pipe 필드에 PIPE 이름을 입력해요.
  6. Apply를 선택해 구성을 저장해요.

오류 처리 커스터마이즈하기

오류 처리는 Openflow 측 실패와 Snowpipe Streaming 서비스 내부의 서버 측 실패로 나뉘어요.

  • Openflow 오류(클라이언트 측 실패): 파싱 불가능한 페이로드나 커스텀 변환 실패 같은 오류는 레코드가 Snowflake에 도달하기 전에 발생해요. 기본적으로 이러한 레코드는 버려져요. Openflow에서 이 오류들을 처리할 수 있어요 — ConsumeKafka 프로세서의 parse failure 관계에서 FlowFiles를 사용해요. 전체 과정은 Kafka as destination for DLQ messages와 공유 가이드 Configuring Dead Letter Queue (DLQ) handling을 참고해요.
  • Snowpipe Streaming 오류(서버 측 실패): Snowflake에 성공적으로 도달했지만 대상 테이블의 스키마와 호환되지 않는 레코드(예: 타입 불일치)의 오류는 Snowflake 인프라가 포착해요. 대상 테이블에서 오류 로깅이 활성화되면(error_logging = true) 이러한 실패한 행은 대상 Error table에 자동으로 수집돼요.

성능 튜닝

더 알아보기 (Learn more)