Protobuf 데이터 타입 수집 구성

Protobuf 데이터 타입 수집 구성

이 페이지에서는 스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)가 기본 JsonTreeReader 대신 StandardProtobufReader를 사용해 Protobuf 인코딩 메시지를 수집하도록 전환하는 방법을 설명해요. 지원되는 스키마 접근 전략과 메시지 이름 해석 방식을 단계별로 다룹니다.

출처: Snowflake 문서

본문

Note

이 커넥터는 Snowflake Connector Terms에 의해 규율됩니다.

스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)는 Consume* 프로세서를 JsonTreeReader 컨트롤러 서비스와 함께 사용해 들어오는 메시지를 파싱합니다. 이 토픽은 JsonTreeReader를 StandardProtobufReader로 교체해 Protobuf 인코딩 메시지로 전환하는 방법을 설명합니다. 단계는 두 커넥터에 동일합니다 — Kinesis 커넥터에서는 ConsumeKafka가 언급된 곳에 ConsumeKinesis 프로세서를 사용하세요.

Tip

이 사용자 정의를 손으로 적용할 필요 없습니다. Snowflake CoCo의 Openflow 스킬이 대신 수행할 수 있습니다 — 원하는 변경을 설명하면 이 페이지의 단계에 따라 흐름을 편집합니다. 구성 요소를 수동으로 구성하는 대신 이 스킬을 사용하는 것을 권장합니다.

Note

AWS Glue Schema Registry는 이 커넥터들의 Protobuf에서는 지원되지 않습니다. 아래의 인라인 스키마 텍스트나 Confluent Schema Registry를 사용하세요.

세 가지 스키마 접근 전략이 지원됩니다:

  • 인라인 스키마 텍스트 (Inline schema text) — 서비스 구성에 Proto 3 스키마를 직접 제공.

  • 스키마 이름 (Schema name) — 구성된 SchemaRegistry 서비스(예: ConfluentSchemaRegistry)로 스키마 이름을 해석.

  • 스키마 참조 리더 (Schema reference reader) — ConfluentEncodedSchemaReferenceReader로 메시지에서 스키마 ID를 읽고 ConfluentSchemaRegistry에 대해 해석. Confluent 와이어 형식 인코딩 메시지의 표준 전략입니다.

추가로, 단일 .proto 파일은 여러 메시지 타입을 정의할 수 있으므로 올바른 메시지 이름을 어떻게 해석할지도 구성해야 합니다:

  • 메시지 이름 속성 (Message name property) — 정규화된 메시지 이름(패키지 포함)을 파라미터로 직접 지정. 예: mypackage.MyMessage.

  • 메시지 이름 리졸버 (Message name resolver) — MessageNameResolver 서비스를 사용해 FlowFile 콘텐츠나 속성에서 메시지 이름을 동적으로 해석. ConfluentProtobufMessageNameResolver는 Confluent 와이어 형식에서 메시지 인덱스를 디코딩하고 스키마 정의에서 정규화된 이름을 조회해 이름을 해석합니다.

전제 조건

  • Openflow에 배포된 기존 Kafka 고성능 또는 Kinesis 고성능 커넥터가 있어야 합니다.

  • Kafka 토픽 / Kinesis 스트림이 Protobuf 인코딩 메시지를 생성해야 합니다.

  • Confluent Schema Registry를 사용한다면: 레지스트리 URL과 Openflow 런타임에서 그 레지스트리로의 네트워크 접근이 있어야 합니다. Snowflake 관리 배포(SPCS)에서는 레지스트리로의 아웃바운드 연결을 허용하도록 적절한 External Access Integration이 구성되고 런타임에 할당되어 있어야 합니다.

Protobuf 데이터 타입 수집 구성

아래 단계는 Kafka 고성능과 Kinesis 고성능 커넥터에 동일하게 적용됩니다(Kinesis 커넥터에서는 ConsumeKafka가 언급된 곳에 ConsumeKinesis 프로세서를 사용).

1단계: StandardProtobufReader 컨트롤러 서비스 생성

  1. Openflow UI에서 커넥터의 프로세스 그룹을 엽니다.

  2. Configure > Controller Services(톱니바퀴 아이콘)로 이동합니다.

  3. **+**를 선택해 새 컨트롤러 서비스를 추가합니다.

  4. StandardProtobufReader(org.apache.nifi.services.protobuf.StandardProtobufReader)를 검색해 선택합니다.

  5. Add를 선택합니다.

2단계: StandardProtobufReader에서 스키마 접근 전략 구성

1단계에서 만든 StandardProtobufReader 컨트롤러 서비스에서 스키마 접근 전략과 메시지 이름 해석 전략을 구성합니다. Protobuf 스키마와 메시지 타입이 분산되는 방식과 일치하는 조합을 선택하세요.

옵션 A — 메시지 이름 속성과 인라인 스키마 텍스트

Proto 3 스키마를 텍스트 문자열로 가지고 있고 메시지 이름이 런타임에 변하지 않을 때 사용하세요.

  1. StandardProtobufReader 서비스의 Edit(톱니바퀴) 아이콘을 선택합니다.

  2. 다음 속성을 설정합니다:

Property Value
Schema Access Strategy Use 'Schema Text' Property
Schema Text The full Proto 3 formatted schema text. Supports Expression Language.
Message Name Resolution Strategy Message Name Property
Message Name The fully qualified name of the Protocol Buffers message including its package, for example, mypackage.MyMessage. Supports Expression Language.
  1. Apply를 선택합니다.

옵션 B — Confluent 메시지 이름 리졸버를 사용한 Confluent Schema Registry

메시지가 Confluent 와이어 형식(매직 바이트 0x00 뒤에 4바이트 스키마 ID)으로 인코딩되고 메시지 이름을 자동 해석해야 할 때 사용하세요. 추가 컨트롤러 서비스 세 개가 필요합니다:

  • ConfluentSchemaRegistry — 스키마 ID를 실제 Proto 3 스키마로 해석.

  • ConfluentEncodedSchemaReferenceReader — 각 메시지에서 스키마 ID를 읽음.

  • ConfluentProtobufMessageNameResolver — Confluent 와이어 형식에서 메시지 인덱스를 디코딩해 정규화된 메시지 이름을 해석하는 메시지 이름 리졸버.

ConfluentSchemaRegistry 컨트롤러 서비스 생성:

  1. Configure > Controller Services로 이동합니다.

  2. **+**를 선택하고 ConfluentSchemaRegistry(org.apache.nifi.confluent.schemaregistry.ConfluentSchemaRegistry)를 검색합니다.

  3. Add를 선택합니다.

  4. Edit(톱니바퀴) 아이콘을 선택하고 다음 속성을 설정합니다:

Property Value
Schema Registry URLs (required) Comma-separated URL(s) of your Confluent Schema Registry, for example, https://schema-registry.example.com:8081.
SSL Context Service (optional) An SSLContextService if the registry requires TLS. Implementations: StandardSSLContextService, StandardRestrictedSSLContextService, PEMEncodedSSLContextProvider.
Communications Timeout (required) How long to wait for a response from the registry before failing.
Cache Size (required) Number of schemas to cache locally. Raise it when a single topic/stream carries many distinct Protobuf message types at high throughput (more memory); the default is fine for a few schemas.
Cache Expiration (required) How long cached schemas are valid before being re-fetched.
Authentication Type NONE or BASIC if the registry requires HTTP Basic authentication.
Username Username for Basic authentication. Only used when Authentication Type is BASIC.
Password Password for Basic authentication. Sensitive property. Only used when Authentication Type is BASIC.
  1. Apply를 선택합니다.

  2. Enable(번개 아이콘)을 선택하고 상태가 Enabled로 표시될 때까지 기다립니다.

ConfluentEncodedSchemaReferenceReader 컨트롤러 서비스 생성:

  1. Configure > Controller Services로 이동합니다.

  2. **+**를 선택하고 ConfluentEncodedSchemaReferenceReader(org.apache.nifi.confluent.schemaregistry.ConfluentEncodedSchemaReferenceReader)를 검색합니다.

  3. Add를 선택합니다.

  4. Enable 아이콘을 선택하고 상태가 Enabled로 표시될 때까지 기다립니다.

Note

ConfluentEncodedSchemaReferenceReader에는 구성 가능한 속성이 없습니다. 각 메시지의 시작 부분에서 Confluent 인코딩 스키마 ID(매직 바이트 0x00 + 4바이트 정수)를 읽을 뿐입니다.

ConfluentProtobufMessageNameResolver 컨트롤러 서비스 생성:

  1. Configure > Controller Services로 이동합니다.

  2. **+**를 선택하고 ConfluentProtobufMessageNameResolver(org.apache.nifi.confluent.schemaregistry.ConfluentProtobufMessageNameResolver)를 검색합니다.

  3. Add를 선택합니다.

  4. Enable 아이콘을 선택하고 상태가 Enabled로 표시될 때까지 기다립니다.

Note

ConfluentProtobufMessageNameResolver에는 구성 가능한 속성이 없습니다. Confluent Protobuf 와이어 형식에 내장된 메시지 인덱스 시퀀스를 디코딩해 정규화된 메시지 이름을 결정합니다. 와이어 형식에 대한 자세한 내용은 Confluent Schema Registry documentation을 참고하세요.

StandardProtobufReader가 스키마 참조 리더와 메시지 이름 리졸버를 사용하도록 구성:

  1. 1단계에서 만든 StandardProtobufReader 서비스의 Edit(톱니바퀴) 아이콘을 선택합니다.

  2. 다음 속성을 설정합니다:

Property Value
Schema Access Strategy Schema Reference Reader
Schema Reference Reader Select the ConfluentEncodedSchemaReferenceReader created above.
Schema Registry Select the ConfluentSchemaRegistry created above.
Message Name Resolution Strategy Message Name Resolver
Message Name Resolver Select the ConfluentProtobufMessageNameResolver created above.
  1. Apply를 선택합니다.

3단계: StandardProtobufReader 컨트롤러 서비스 활성화

StandardProtobufReader의 Enable(번개 아이콘)을 선택하고 상태가 Enabled로 표시될 때까지 기다립니다.

4단계: 소스 프로세서 업데이트

  1. ConsumeKafka 프로세서(Kinesis 커넥터에서는 ConsumeKinesis)를 더블클릭해 속성을 엽니다.

  2. 다음 속성을 업데이트합니다:

Property Value
Record Reader Select the StandardProtobufReader created above.
  1. Apply를 선택합니다.

5단계: JsonTreeReader 컨트롤러 서비스 비활성화

원래 JsonTreeReader는 더 이상 필요하지 않습니다.

  1. 프로세스 그룹이 실행 중이면 중지합니다.

  2. Configure > Controller Services로 이동합니다.

  3. JsonTreeReader의 Disable 아이콘을 선택합니다.

  4. 서비스가 다른 프로세서에서 참조하지 않는다면 Delete(휴지통) 아이콘을 선택해 삭제할 수도 있습니다.

  5. 프로세스 그룹을 시작합니다.

문제 해결

Symptom Likely cause
SchemaNotFoundException at runtime The schema ID in the message is not present in the registry, or the registry URL is misconfigured. Verify Schema Registry URLs in ConfluentSchemaRegistry.
Parse failures or malformed records The Schema Text does not match the actual message schema, or the wrong Message Name is specified. Compare the schema and message name with those used by the producer. Failed messages are routed to the parse-failure relationship (parse failure on Kafka, parse.failure on Kinesis).
ConfluentSchemaRegistry fails to enable Network connectivity issue between the Openflow runtime and the registry. Check the External Access Integration and that the registry URL is reachable from the cluster.
Authentication failure against registry Registry requires Basic auth — set Authentication Type to BASIC and provide Username and Password in ConfluentSchemaRegistry.

더 알아보기 (Learn more)