Amazon Kinesis Data Streams용 Openflow 커넥터 설정
Amazon Kinesis Data Streams용 Openflow 커넥터 설정
이 문서에서는 Amazon Kinesis Data Streams용 Openflow 커넥터를 설정하는 방법을 설명합니다. 이 커넥터는 스키마 진화(schema evolution) 기능을 갖추고 Kinesis 스트림의 JSON 메시지를 Snowflake 테이블로 수집하도록 설계되었습니다.
출처: Snowflake 문서
본문
Kinesis용 Openflow 커넥터 설정
사전 준비 사항
- Amazon Kinesis Data Streams용 Openflow 커넥터 문서를 검토하세요.
- Set up Openflow - BYOC 또는 Set up Openflow - Snowflake Deployments를 완료했는지 확인하세요.
- Openflow - Snowflake Deployments를 사용한다면 필수 도메인(required domains) 구성 문서를 검토하고, Kinesis 커넥터에 필요한 필수 도메인에 대한 접근 권한을 부여했는지 확인하세요.
AWS에 IAM 역할 및 정책 설정
AWS 관리자는 AWS 계정에서 다음 작업을 수행하세요.
- Openflow가 Kinesis 데이터 스트림에 접근할 때 사용할 AWS IAM 사용자 또는 역할을 만듭니다. 자세한 내용은 AWS 문서의 'Creating IAM users'를 참조하세요.
- AWS 사용자에게 Access Key 자격 증명이 구성되어 있는지 확인하세요.
- AWS 사용자에게 다음 IAM 권한을 부여하세요.
| 서비스 | 작업(Actions) | 리소스(ARNs) | 목적 |
|---|---|---|---|
| Amazon Kinesis Data Streams | kinesis:DescribeStream, kinesis:DescribeStreamConsumer, kinesis:GetRecords, kinesis:GetShardIterator, kinesis:ListShards, kinesis:RegisterStreamConsumer | arn:aws:kinesis:${REGION}:${ACCOUNT_ID}:stream/${STREAM_NAME} | 샤드를 발견하고, 공유 처리량(shred-throughput) 폴링으로 레코드를 읽고, 스트림 ARN을 확인하고, Enhanced Fan-Out 소비자를 등록하며, 등록 중 소비자 상태를 폴링합니다. |
| Amazon Kinesis Data Streams | kinesis:DeregisterStreamConsumer, kinesis:DescribeStreamConsumer, kinesis:SubscribeToShard | arn:aws:kinesis:${REGION}:${ACCOUNT_ID}:stream/${STREAM_NAME}/consumer/* | 소비자 ARN으로 Enhanced Fan-Out 소비자를 설명하고 구독하며 등록을 해제합니다. |
| Amazon DynamoDB | dynamodb:CreateTable, dynamodb:DeleteTable, dynamodb:DescribeTable, dynamodb:GetItem, dynamodb:PutItem, dynamodb:Query, dynamodb:Scan, dynamodb:UpdateItem | arn:aws:dynamodb:${REGION}:${ACCOUNT_ID}:table/${APPLICATION_NAME}, arn:aws:dynamodb:${REGION}:${ACCOUNT_ID}:table/${APPLICATION_NAME}_migration | 체크포인트/임대 테이블(샤드 임대, 노드 하트비트, 체크포인트)과, 레거시 체크포인트 테이블에서 1회성 마이그레이션 중 사용되는 임시 마이그레이션 테이블을 만들고 관리합니다. |
예시 IAM 정책:
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "KinesisStreamAccess",
"Effect": "Allow",
"Action": [
"kinesis:DescribeStream",
"kinesis:DescribeStreamConsumer",
"kinesis:GetRecords",
"kinesis:GetShardIterator",
"kinesis:ListShards",
"kinesis:RegisterStreamConsumer"
],
"Resource": "arn:aws:kinesis:${REGION}:${ACCOUNT_ID}:stream/${STREAM_NAME}"
},
{
"Sid": "KinesisConsumerAccess",
"Effect": "Allow",
"Action": [
"kinesis:DeregisterStreamConsumer",
"kinesis:DescribeStreamConsumer",
"kinesis:SubscribeToShard"
],
"Resource": "arn:aws:kinesis:${REGION}:${ACCOUNT_ID}:stream/${STREAM_NAME}/consumer/*"
},
{
"Sid": "DynamoDBTableAccess",
"Effect": "Allow",
"Action": [
"dynamodb:CreateTable",
"dynamodb:DeleteTable",
"dynamodb:DescribeTable",
"dynamodb:GetItem",
"dynamodb:PutItem",
"dynamodb:Query",
"dynamodb:Scan",
"dynamodb:UpdateItem"
],
"Resource": [
"arn:aws:dynamodb:${REGION}:${ACCOUNT_ID}:table/${APPLICATION_NAME}",
"arn:aws:dynamodb:${REGION}:${ACCOUNT_ID}:table/${APPLICATION_NAME}_migration"
]
}
]
}
예시 정책을 사용하기 전에 다음 플레이스홀더를 바꿔주세요.
| 플레이스홀더 | 설명 |
|---|---|
| ${REGION} | AWS 리전(예: us-east-1) |
| ${ACCOUNT_ID} | AWS 계정 ID(예: 123456789012) |
| ${STREAM_NAME} | AWS Kinesis Stream Name 커넥터 파라미터 값 |
| ${APPLICATION_NAME} | AWS Kinesis Application Name 커넥터 파라미터 값. DynamoDB 체크포인트 테이블 이름과 Enhanced Fan-Out 등록 소비자 이름으로 사용됩니다. |
참고:
${APPLICATION_NAME}_migration테이블은 레거시 체크포인트 테이블에서 새 스키마로의 1회성 마이그레이션 동안에만 생성되는 임시 DynamoDB 테이블입니다. 마이그레이션이 완료되면 자동으로 삭제됩니다. 레거시 KCL 기반 커넥터를 사용한 적이 없다면 정책에서 마이그레이션 테이블 ARN을 생략할 수 있습니다.dynamodb:DeleteTable액션은 마이그레이션 과정에서 사용되며, 마이그레이션이 완료된 것이 확인된 후 정책에서 제거할 수 있습니다.kinesis:DeregisterStreamConsumer액션은 프로세서가 캔버스에서 제거될 때 호출됩니다. IAM 주체에 이 권한이 없다면 AWS 콘솔이나 CLI를 통해 소비자를 수동으로 등록 해제해야 합니다.
Snowflake 계정 설정
Snowflake 계정 관리자는 다음 작업을 수행하세요.
- 유형이 SERVICE인 새 Snowflake 서비스 사용자를 만듭니다.
- 새 역할을 만들거나 기존 역할을 사용해 데이터베이스 권한을 부여합니다. 커넥터는 목적지 테이블을 만들어야 하므로, 사용자가 Snowflake 객체를 관리하는 데 필요한 권한을 보유하는지 확인하세요.
| 객체 | 권한 | 참고 |
|---|---|---|
| Database | USAGE | |
| Schema | USAGE | |
| Table | OWNERSHIP | 커넥터가 테이블에 데이터를 수집하는 데 필요합니다. |
Snowflake는 더 나은 접근 제어를 위해 Kinesis 스트림마다 별도의 사용자와 역할을 만들 것을 권장합니다. 다음 스크립트로 사용자 지정 역할을 만들고 구성할 수 있습니다( SECURITYADMIN 또는 이에 상응하는 역할 필요):
USE ROLE securityadmin;
CREATE ROLE openflow_kinesis_connector_role_1;
GRANT USAGE ON DATABASE kinesis_db TO ROLE openflow_kinesis_connector_role_1;
GRANT USAGE ON SCHEMA kinesis_schema TO ROLE openflow_kinesis_connector_role_1;
참고: 권한은 커넥터 역할에 직접 부여해야 하며 상속될 수 없습니다.
- 목적지 테이블 구성: Snowflake는 스키마 변경에는 서버 측 스키마 진화(server-side schema evolution)를, DML 오류 로깅에는 오류 테이블을 권장합니다. 아래 예시는 테이블을 만들고 OWNERSHIP 권한을 추가하는 방법을 보여줍니다.
USE ROLE openflow_kinesis_connector_role_1;
CREATE TABLE kinesis_db.kinesis_schema.<DESTINATION_TABLE_NAME> (
kinesisMetadata object
)
ENABLE_SCHEMA_EVOLUTION = TRUE
ERROR_LOGGING = TRUE;
USE ROLE securityadmin;
GRANT OWNERSHIP ON TABLE <DESTINATION_TABLE_NAME> TO ROLE openflow_kinesis_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 캔버스의 오른쪽 위 햄버거 메뉴에서 이 Secrets Manager와 연결된 Parameter Provider를 구성하세요. Controller Settings » Parameter Provider로 이동해 파라미터 값을 가져오세요.
- 이 시점부터 모든 자격 증명을 연결된 파라미터 경로로 참조할 수 있으며, Openflow 내에 민감한 값을 영구 저장할 필요가 없습니다.
- 사용자에게 접근 권한 부여: 커넥터가 수집한 원시 데이터에 접근해야 하는 다른 Snowflake 사용자(예: Snowflake에서 사용자 지정 처리를 위해)에게는 2단계에서 만든 역할을 부여해야 합니다.
(선택) 아웃바운드 AWS PrivateLink 구성
Openflow - Snowflake Deployments에서 커넥터를 실행 중이고 공개 인터넷 대신 아웃바운드 프라이빗 연결(AWS PrivateLink)로 커넥터의 Kinesis 트래픽을 라우팅하고 싶다면 이 섹션의 단계를 따르세요.
커넥터는 다음 AWS 서비스에 아웃바운드 호출을 합니다.
| 서비스 | 목적 | PrivateLink 지원 |
|---|---|---|
| Amazon Kinesis Data Streams | 스트림 레코드 읽기 | 이 커넥터에서 지원함 |
| Amazon DynamoDB | 처리된 레코드의 체크포인트 메타데이터 저장 | 미지원, 공개 엔드포인트 사용 |
중요: Amazon DynamoDB는 PrivateLink 엔드포인트에 대해 Private DNS를 지원하지 않습니다. AWS 문서의 'Considerations when using AWS PrivateLink for Amazon DynamoDB'를 참조하세요. Snowflake의
PRIVATE_HOST_PORT네트워크 규칙 유형은 Private DNS에 의존하므로, 커넥터는 DynamoDB 트래픽을 PrivateLink 엔드포인트로 라우팅할 수 없습니다. 예시와 같이HOST_PORT(공개 엔드포인트)를 사용해 DynamoDB를 구성하세요. DynamoDB로는 체크포인트 메타데이터만 공개 엔드포인트를 통해 흐릅니다. 스트림 레코드는 프라이빗 Kinesis 엔드포인트를 통해 흐릅니다.
아웃바운드 AWS PrivateLink를 구성하려면 다음 단계를 완료하세요.
- ACCOUNTADMIN으로 스트림이 있는 리전에 Amazon Kinesis Data Streams용 아웃바운드 PrivateLink 엔드포인트를 프로비저닝합니다.
을 AWS 리전(예: us-east-1)으로 바꾸세요.
USE ROLE ACCOUNTADMIN;
SELECT SYSTEM$PROVISION_PRIVATELINK_ENDPOINT(
'com.amazonaws.<region>.kinesis-streams',
'kinesis.<region>.amazonaws.com'
);
자세한 내용은 SYSTEM$PROVISION_PRIVATELINK_ENDPOINT와 'Managing outbound private connectivity endpoints on AWS'를 참조하세요.
- 프라이빗 엔드포인트로 kinesis.
.amazonaws.com에 도달하고 DynamoDB는 공개 엔드포인트로 도달하는 네트워크 규칙을 만듭니다. <openflow_network_schema>를 네트워크 규칙을 호스팅하는 스키마로 바꾸세요.
USE ROLE ACCOUNTADMIN;
USE SCHEMA <openflow_network_schema>;
CREATE OR REPLACE NETWORK RULE openflow_kinesis_private_network_rule
MODE = EGRESS
TYPE = PRIVATE_HOST_PORT
VALUE_LIST = ('kinesis.<region>.amazonaws.com');
CREATE OR REPLACE NETWORK RULE openflow_kinesis_public_network_rule
MODE = EGRESS
TYPE = HOST_PORT
VALUE_LIST = ('dynamodb.<region>.amazonaws.com:443');
- 두 네트워크 규칙을 외부 접근 통합(external access integration)에 연결한 뒤, 실행 역할(execute-as role)에 통합 사용 권한을 부여합니다.
USE ROLE ACCOUNTADMIN;
CREATE OR REPLACE EXTERNAL ACCESS INTEGRATION openflow_kinesis_eai
ALLOWED_NETWORK_RULES = (
openflow_kinesis_private_network_rule,
openflow_kinesis_public_network_rule
)
ENABLED = TRUE
COMMENT = 'External access integration for the Openflow Connector for Kinesis';
GRANT USAGE ON INTEGRATION openflow_kinesis_eai TO ROLE OPENFLOW_<RUNTIME_NAME>_EXECUTE_AS_RL;
통합을 런타임과 연결하는 단계는 'Set up Openflow - Snowflake Deployment: Configure allowed domains for Openflow connectors'를 참조하세요.
커넥터 설정
데이터 엔지니어는 커넥터를 설치하고 구성하려면 다음 작업을 수행하세요.
커넥터 설치
- Openflow의 Connector library 탭으로 이동하세요.
- Openflow connectors 페이지에서 Amazon Kinesis Data Streams용 Openflow 커넥터를 찾아 Install을 선택하세요.
- Select runtime 대화상자의 Available runtimes 드롭다운 목록에서 런타임을 선택하고 Add를 클릭하세요.
- 참고: 커넥터를 설치하기 전에 Snowflake에 커넥터가 수집한 데이터를 저장할 데이터베이스, 스키마, 테이블을 만들었는지 확인하세요.
- Snowflake 계정 자격 증명으로 배포에 인증하고, 런타임 애플리케이션이 Snowflake 계정에 접근하도록 허용할지 묻는 프롬프트에서 Allow를 선택하세요. 커넥터 설치 과정은 완료까지 몇 분 걸립니다.
- Snowflake 계정 자격 증명으로 런타임에 인증하세요. 커넥터 프로세스 그룹이 추가된 Openflow 캔버스가 표시됩니다.
커넥터 구성
- 필요하면 내장 파라미터를 구성하기 전에 커넥터 구성을 사용자 지정하세요.
- 프로세스 그룹 파라미터 채우기: 가져온 프로세스 그룹을 마우스 오른쪽 버튼으로 클릭하고 Parameters를 선택하세요.
- 필수 파라미터 값을 채우세요.
일반 파라미터:
| 파라미터 | 설명 | 필수 |
|---|---|---|
| AWS Access Key ID | Kinesis Stream과 DynamoDB에 연결하기 위한 AWS Access Key ID | 예 |
| AWS Kinesis Region | 연결할 AWS 리전. 일반 AWS 리전 형식을 사용하세요(예: us-west-2, ap-southeast-1, eu-west-1). AWS Regions 페이지 참조 | 예 |
| AWS Secret Access Key | Kinesis Stream과 DynamoDB에 연결하기 위한 AWS Secret Access Key | 예 |
| AWS Kinesis Application Name | Kinesis Stream 소비 진행 상황 추적을 위한 DynamoDB 테이블 이름으로 사용되는 이름 | 예 |
| AWS Kinesis Consumer Type | Kinesis Stream에서 레코드를 읽는 전략. 다음 값 중 하나여야 함: SHARED_THROUGHPUT, ENHANCED_FAN_OUT. 자세한 내용은 'Differences between shared throughput consumer and enhanced fan-out consumer' 참조 | 예 |
| AWS Kinesis Initial Stream Position | 데이터 복제가 시작되는 초기 스트림 위치. 특정 AWS Kinesis Application Name의 최초 시작 시에만 적용됨. 가능한 값: LATEST(가장 최근 저장 레코드), TRIM_HORIZON(가장 오래된 저장 레코드) | 예 |
| AWS Kinesis Stream Name | 데이터를 소비할 AWS Kinesis Stream Name | 예 |
| 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에 이미 존재해야 함. 이름은 대소문자를 구분함. 따옴표 없는 식별자는 대문자로 입력 | 예 |
커넥터 시작
- 캔버스를 마우스 오른쪽 버튼으로 클릭하고 Enable all Controller Services를 선택하세요.
- 캔버스를 마우스 오른쪽 버튼으로 클릭하고 Start를 선택하세요. 커넥터가 데이터 수집을 시작합니다.
KINESISMETADATA 열 이해
커넥터는 KINESISMETADATA 구조에 Kinesis 레코드에 대한 메타데이터를 채웁니다. 구조에는 다음 정보가 포함됩니다.
| 필드 이름 | 필드 유형 | 예시 값 | 설명 |
|---|---|---|---|
| stream | String | stream-name | 레코드가 온 Kinesis 스트림 이름 |
| shardId | String | shardId-000000000001 | 레코드가 온 스트림의 샤드 식별자 |
| approximateArrival | Number | 1782092074893 | 레코드가 스트림에 삽입된 대략적 시간(Unix epoch 밀리초 타임스탬프) |
| partitionKey | String | key-1234 | 데이터 생성자가 레코드에 지정한 파티션 키 |
| sequenceNumber | String | 123456789 | Kinesis Data Streams가 샤드의 레코드에 할당한 고유 시퀀스 번호 |
| subSequenceNumber | Number | 2 | 레코드의 하위 시퀀스 번호(동일 시퀀스 번호를 가진 집계 레코드에 사용됨) |
| shardedSequenceNumber | String | 12345678900002 | 레코드의 시퀀스 번호와 하위 시퀀스 번호의 조합 |
수집 지연 시간 측정
변경 추적, 증분 처리, 행 수정 시간 기반 Time Travel 쿼리를 위해 ROW_TIMESTAMP 기능을 사용할 수 있습니다.
목적지 테이블에서 다음 명령을 실행해 활성화할 수 있습니다.
ALTER TABLE <DESTINATION_TABLE> SET ROW_TIMESTAMP = TRUE;
행 타임스탬프가 활성화되면 테이블은 각 행이 마지막으로 수정된 시각을 반환하는 METADATA$ROW_LAST_COMMIT_TIME 열을 노출합니다. 자세한 내용은 Row timestamps를 참조하세요.
참고: 행 타임스탬프는 인터랙티브 테이블에서는 사용할 수 없습니다. 자세한 내용은 'Limitations of interactive tables'를 참조하세요.
Apache Iceberg™ 테이블과 함께 커넥터 사용
커넥터는 Snowflake 관리형 Apache Iceberg™ 테이블에 데이터를 수집할 수 있습니다. 커넥터는 Iceberg 테이블을 자동으로 만들지 않으므로, 커넥터를 실행하기 전에 Iceberg 테이블을 수동으로 만들어야 합니다.
커넥터는 표준 Snowflake 테이블과 같은 방식으로 Iceberg 목적지 테이블에 대해 서버 측 스키마 진화를 지원합니다. 목적지 테이블에 ENABLE_SCHEMA_EVOLUTION = TRUE가 있으면 Snowflake가 수신 스트림에서 감지된 새 열을 자동으로 추가하고, 새 데이터 패턴을 수용하도록 NOT NULL 제약 조건을 제거합니다. 스키마 진화 동작에 대한 자세한 내용은 Table schema evolution을 참조하세요.
Iceberg 테이블은 다음 스토리지 옵션 중 하나를 사용할 수 있습니다.
- Snowflake 스토리지: Snowflake가 Iceberg 테이블 파일을 저장하고 관리해 주므로 외부 볼륨을 만들거나 커넥터에 접근 권한을 부여할 필요가 없습니다.
- 직접 관리하는 외부 클라우드 스토리지(외부 볼륨을 통해 접근). 커넥터 역할에 외부 볼륨에 대한 USAGE를 부여해야 합니다.
외부 볼륨에 사용 권한 부여
이 단계는 Iceberg 테이블이 직접 관리하는 외부 볼륨을 사용할 때만 적용됩니다. 테이블이 Snowflake 스토리지를 사용한다면 이 단계를 건너뛰세요.
예를 들어 Iceberg 테이블이 kinesis_external_volume 외부 볼륨을 사용하고 커넥터가 openflow_kinesis_connector_role_1 역할을 사용한다면 다음 문을 실행하세요.
USE ROLE ACCOUNTADMIN;
GRANT USAGE ON EXTERNAL VOLUME kinesis_external_volume TO ROLE openflow_kinesis_connector_role_1;
수집용 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 스토리지를 사용하려면 EXTERNAL_VOLUME = 'SNOWFLAKE_MANAGED'로 설정하고 BASE_LOCATION은 생략하세요.
CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
kinesisMetadata OBJECT(
stream STRING,
shardId STRING,
approximateArrival STRING,
partitionKey STRING,
sequenceNumber STRING,
subSequenceNumber INTEGER,
shardedSequenceNumber STRING
),
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 (
kinesisMetadata OBJECT(
stream STRING,
shardId STRING,
approximateArrival BIGINT,
partitionKey STRING,
sequenceNumber STRING,
subSequenceNumber INTEGER,
shardedSequenceNumber STRING
),
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)은 낮은 지연 시간과 높은 동시성을 위해 최적화된 특수한 유형의 Snowflake 테이블입니다. 인터랙티브 테이블에 대한 자세한 내용은 인터랙티브 테이블 문서에서 확인할 수 있습니다.
- 인터랙티브 테이블을 만듭니다.
CREATE INTERACTIVE TABLE REALTIME_METRICS (
metric_name VARCHAR,
metric_value NUMBER,
source_stream VARCHAR,
approximate_arrival NUMBER
) CLUSTER BY (metric_name)
AS (SELECT
$1:M_NAME::VARCHAR,
$1:M_VALUE::NUMBER,
$1:kinesisMetadata.stream::VARCHAR,
$1:kinesisMetadata.approximateArrival::NUMBER
from TABLE(DATA_SOURCE(TYPE => 'STREAMING')));
중요한 고려사항:
- 인터랙티브 테이블에는 특정한 제한과 쿼리 제약이 있습니다. 커넥터와 함께 사용하기 전에 인터랙티브 테이블 문서를 검토하세요.
- 인터랙티브 테이블의 경우 필요한 모든 변환은 테이블 정의에서 처리해야 합니다.
- 인터랙티브 테이블을 효율적으로 쿼리하려면 인터랙티브 웨어하우스가 필요합니다.
목적지 테이블의 사용자 지정 스키마와 함께 커넥터 사용
커넥터는 각 Kinesis 레코드를 Snowflake 테이블에 삽입할 행으로 취급합니다. 예를 들어 다음 JSON처럼 구조화된 메시지 콘텐츠를 가진 Kinesis 스트림이 있다면:
{
"order_id": 12345,
"customer_name": "John",
"order_total": 100.00,
"isPaid": true
}
기본적으로 JSON의 모든 필드를 지정할 필요는 없습니다. 스키마 진화가 처리해 줍니다. 하지만 정적 스키마를 선호한다면 다음을 실행해 만들 수 있습니다.
CREATE TABLE ORDERS (
kinesisMetadata 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
FROM (
SELECT
$1:order_id::STRING,
$1:customer_name,
$1:order_total::STRING,
$1:isPaid::STRING
FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING'))
);
자체 파이프를 정의하면 목적지 테이블 열이 JSON 키와 일치할 필요가 없습니다. 원하는 이름으로 열을 바꾸고 필요하면 데이터 유형을 캐스트할 수 있습니다.
커넥터가 사용자 지정 파이프와 함께 동작하도록 조정하려면 다음 작업을 수행하세요.
- Openflow 캔버스에서 Kinesis 수집 흐름에 사용된 PublishSnowpipeStreaming 프로세서를 마우스 오른쪽 버튼으로 클릭하세요.
- 컨텍스트 메뉴에서 Configure를 선택하세요.
- Properties 탭으로 이동하세요.
- Destination type 필드에서 Pipe를 선택하세요.
- Pipe 필드에 파이프 이름을 입력하세요.
- Apply를 선택해 구성을 저장하세요.
오류 처리 사용자 지정
오류 처리는 Openflow 측 실패와 Snowpipe Streaming 서비스 내의 서버 측 실패로 나뉩니다.
- Openflow 오류(클라이언트 측 실패): 파싱할 수 없는 페이로드나 사용자 지정 변환 실패 같은 오류는 레코드가 Snowflake에 도달하기 전에 발생합니다. 기본적으로 이러한 레코드는 폐기됩니다. Openflow에서 이러한 오류를 처리할 수 있습니다. ConsumeKinesis 프로세서의 parse failure 관계에서 FlowFile을 사용하세요. 전체 워크스루는 'Kinesis as destination for DLQ messages'와 공유 문서 'Configuring Dead Letter Queue (DLQ) handling'을 참조하세요.
- Snowpipe Streaming 오류(서버 측 실패): Snowflake에 성공적으로 도달했지만 목적지 테이블 스키마와 호환되지 않는(예: 유형 불일치) 레코드에 대한 오류는 Snowflake 인프라가 캡처합니다. 목적지 테이블에서 오류 로깅이 활성화되면(error_logging = true) 이러한 실패한 행은 자동으로 목적지 오류 테이블에 수집됩니다.
다음 단계
- Amazon Kinesis Data Streams용 Openflow 커넥터 성능 튜닝
- Amazon Kinesis Data Streams용 Openflow 커넥터 유지 관리
- Amazon Kinesis Data Streams용 Openflow 커넥터 문제 해결
- Kinesis 데이터 스트림용 Openflow 커넥터: DLQ 처리 구성
- 사용자 지정 변환 구성(Configuring custom transformations)