Kafka용 Openflow 커넥터 설정
Kafka용 Openflow 커넥터 설정
이 주제는 Openflow Connector for Kafka를 설정하는 단계를 설명해요.
출처: Snowflake 문서
본문
참고: 이 커넥터는 Snowflake Connector Terms가 적용돼요.
전제 조건(Prerequisites)
- Snowflake Openflow Connector for Kafka을 검토했는지 확인해요.
- Set up Openflow - BYOC 또는 Set up Openflow - Snowflake Deployments를 설정했는지 확인해요.
- Openflow - Snowflake Deployments를 사용한다면 configuring required domains을 검토하고, Kafka 커넥터의 필수 도메인에 대한 액세스를 부여했는지 확인해요. 커넥터는 클러스터의 모든 Kafka 브로커에 연결할 수 있어야 해요.
Snowflake 계정 설정하기
Snowflake 계정 관리자로서 다음 작업을 수행해요:
-
타입이 SERVICE인 새 Snowflake 서비스 사용자를 만들어요.
-
새 역할을 만들거나 기존 역할을 사용하고 데이터베이스 권한(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;참고: 권한은 커넥터 역할에 직접 부여되어야 하며 상속될 수 없어요.
-
대상 테이블을 구성해요. 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의 키가 테이블 컬럼과 일치하지 않으면 커넥터는 그 키를 무시해요. -
(선택) 시크릿 매니저를 구성해요. Snowflake는 이 단계를 강력히 권장해요. Openflow가 지원하는 시크릿 매니저(예: AWS, Azure, Hashicorp)를 구성하고 공개·개인 키를 시크릿 저장소에 보관해요.
- 구성 후 시크릿 매니저에 인증할 방법을 결정해요. AWS에서는 Openflow와 연결된 EC2 인스턴스 역할을 사용해 다른 시크릿을 저장할 필요가 없도록 하는 것을 권장해요.
- Openflow에서 오른쪽 상단의 햄버거 메뉴에서 이 시크릿 매니저와 연결된 Parameter Provider를 구성해요. Controller Settings > Parameter Provider로 이동해 파라미터 값을 가져와요.
- 민감한 값이 Openflow 안에 저장되지 않도록 모든 자격 증명을 연결된 파라미터 경로로 참조해요.
사용자에게 액세스 부여 커넥터가 수집한 원시 데이터에 액세스해야 하는 다른 Snowflake 사용자(예: Snowflake에서 커스텀 처리용)에게 1단계에서 만든 역할을 부여해요.
커넥터 설정하기
데이터 엔지니어로서 커넥터를 설치하고 구성하려면 다음 작업을 수행해요:
커넥터 설치하기
커넥터를 설치하려면:
- Openflow의 Connector library 탭으로 이동해요.
- Openflow 커넥터 페이지에서 커넥터를 찾고 Install을 선택해요.
- Select runtime 대화상자에서 Available runtimes 드롭다운 목록에서 런타임을 선택하고 Add를 선택해요.
참고: 커넥터를 설치하기 전에, 수집된 데이터를 저장할 Snowflake에 데이터베이스, 스키마, 테이블을 만들었는지 확인해요.
- Snowflake 계정 자격 증명으로 배포에 인증하고, 런타임 애플리케이션이 Snowflake 계정에 접근하는 것을 허용하라는 프롬프트가 나오면 Allow를 선택해요. 커넥터 설치 과정은 완료되는 데 몇 분이 걸려요.
- Snowflake 계정 자격 증명으로 런타임에 인증해요. Openflow 캔버스에 커넥터 프로세스 그룹이 추가된 상태로 나타나요.
커넥터 구성하기
- 필요하면 내장 파라미터를 구성하기 전에 커넥터 구성을 커스터마이즈해요. 일부 일반적인 커스터마이즈에는 전용 가이드가 있어요(예: custom transformations, Avro 및 Protobuf 데이터 타입 수집, dead-letter-queue 처리). Snowflake CoCo의 Openflow 스킬로 커스터마이즈를 적용할 수도 있어요. 자세한 내용은 Configuring custom transformations을 참고해요.
- 프로세스 그룹 파라미터를 채워요.
- 가져온 프로세스 그룹을 마우스 오른쪽 버튼으로 클릭하고 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에 이미 존재해야 함. 이름은 대소문자를 구분함. 인용되지 않은 식별자는 대문자로 제공 | 예 |
커넥터 시작하기
- 평면(plane)을 마우스 오른쪽 버튼으로 클릭하고 Enable all Controller Services를 선택해요.
- 평면을 마우스 오른쪽 버튼으로 클릭하고 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을 참고해요.
- 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 키와 일치할 필요가 없어요. 컬럼을 원하는 이름으로 바꾸고 필요하면 데이터 타입을 캐스트할 수 있어요.
커넥터가 커스텀 파이프로 작동하도록 조정하려면:
- Openflow 캔버스에서 Kafka 수집 플로우에 사용된 PublishSnowpipeStreaming 프로세서를 마우스 오른쪽 버튼으로 클릭해요.
- 컨텍스트 메뉴에서 Configure를 선택해요.
- Properties 탭으로 이동해요.
- Destination type 필드에서 Pipe를 선택해요.
- Pipe 필드에 PIPE 이름을 입력해요.
- 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에 자동으로 수집돼요.