ConsumeKafka 프로세서

ConsumeKafka 프로세서

이 문서는 Snowflake OpenFlow의 ConsumeKafka 프로세서에 대한 참조 문서예요. Apache Kafka Consumer API를 사용해 Kafka 메시지를 소비하고, 각 FlowFile에 하나 이상의 직렬화된 Kafka 레코드를 담아 보내요.

출처: Snowflake 문서

본문

기능 — 일반 제공 (Generally Available)

Openflow Snowflake 배포는 AWS, Azure, GCP Commercial 리전의 모든 계정에서 사용할 수 있어요.

Openflow BYOC 배포는 AWS Commercial 리전의 모든 계정에서 사용할 수 있어요.

번들 (Bundle)

그룹 NAR
com.snowflake.openflow.runtime runtime-kafka-nar

설명 (Description)

Apache Kafka Consumer API에서 메시지를 소비해요. 메시지 전송을 위한 보완 NiFi 프로세서는 PublishKafka예요. 이 프로세서는 Kafka 메시지 소비를 지원하며, 선택적으로 NiFi 레코드로 해석할 수 있어요. 현재(레코드 읽기 모드에서) 이 프로세서는 주어진 파티션에서 가져온 모든 레코드가 동일한 스키마를 갖는다고 가정한다는 점에 유의해야 해요. 이 모드에서 Kafka 메시지 중 일부를 가져왔지만 구성된 Record Reader 또는 Record Writer로 파싱하거나 기록할 수 없으면 메시지 내용이 별도의 FlowFile에 기록되고 그 FlowFile이 'parse.failure' 관계로 전달돼요. 그 외에는 각 FlowFile이 'success' 관계로 전달되며 단일 FlowFile 안에 많은 개별 메시지를 담을 수 있어요. FlowFile에 포함된 메시지 수를 나타내기 위해 'record.count' 속성이 추가돼요. 두 Kafka 메시지가 서로 다른 스키마를 가지거나, 특정 속성에 포함된 메시지 헤더 값이 서로 다르면 같은 FlowFile에 넣어지지 않아요.

태그 (Tags)

avro, consume, csv, get, ingest, ingress, json, kafka, openflow, pubsub, record, topic

입력 요구사항 (Input Requirement)

FORBIDDEN

민감한 동적 속성 지원 (Supports Sensitive Dynamic Properties)

아니요 (false)

속성 (Properties)

속성 설명
Commit Offsets 메시지를 받은 후 이 프로세서가 오프셋을 Kafka에 커밋할지 여부를 지정해요. 일반적으로 수신된 메시지가 중복되지 않도록 이 값을 true로 설정해야 해요. 하지만 특정 시나리오에서는 데이터를 처리한 뒤 PublishKafka가 나중에 승인할 수 있도록 오프셋 커밋을 피하고 싶을 수 있으며, 이렇게 하면 Exactly Once 의미론을 제공할 수 있어요.
Content Field 레코드의 어느 필드 아래에 콘텐츠가 추가될지 지정해요. 설정하지 않으면 콘텐츠가 레코드의 루트에 위치해요.
Group ID Kafka group.id 속성에 해당하는 Kafka Consumer Group 식별자예요.
Header Encoding Kafka Record Header 값을 읽고 FlowFile 속성을 쓸 때 적용되는 문자 인코딩이에요.
Header Name Pattern FlowFile 속성으로 기록할 Header Values를 선택하기 위해 Kafka Record Header Names에 적용되는 정규 표현식 패턴이에요.
Headers Field Parent 레코드의 어느 필드 아래에 headers 필드가 추가될지 지정해요. 설정하지 않으면 headers 필드가 레코드의 루트에 위치해요.
Kafka Connection Service Kafka Records를 게시하기 위해 Kafka Broker에 대한 연결을 제공해요.
Key Attribute Encoding Kafka Record Key를 포함하는 구성된 FlowFile 속성 값의 인코딩이에요.
Key Field Parent 레코드의 어느 필드 아래에 key 필드가 추가될지 지정해요. 설정하지 않으면 key 필드가 레코드의 루트에 위치해요.
Key Format Kafka Record Key를 출력 FlowFile에서 어떻게 표현할지 지정해요.
Key Record Reader Kafka Record Key를 Record로 파싱하는 데 사용할 Record Reader예요.
Max Uncommitted Time 프로세서가 FlowFiles를 흐름을 따라 전달하고 (적절하다면) 오프셋을 Kafka에 커밋하기 전에 Kafka에서 소비할 수 있는 최대 시간을 지정해요. 시간이 길수록 지연 시간이 늘어날 수 있어요.
Message Demarcator KafkaConsumer는 메시지를 배치로 받으므로 이 프로세서에는 주어진 토픽과 파티션에 대해 단일 배치의 모든 Kafka 메시지를 담은 FlowFiles를 출력하는 옵션이 있어요. 이 속성을 사용하면 여러 Kafka 메시지를 구분하는 데 사용할 문자열(UTF-8로 해석)을 제공할 수 있어요. 선택 속성이며 제공하지 않으면 수신된 각 Kafka 메시지가 트리거될 때 단일 FlowFile이 돼요. 'new line' 같은 특수 문자를 입력하려면 OS에 따라 CTRL+Enter 또는 Shift+Enter를 사용해요.
Metadata Field 레코드의 어느 필드 아래에 메타데이터가 추가될지 지정해요. 설정하지 않으면 메타데이터가 레코드의 루트에 위치해요.
Metadata Received Timestamp Field 지정하면 출력 FlowFile의 레코드 메타데이터에서 지정된 필드 아래에 타임스탬프가 배치돼요.
Output Strategy Kafka Record를 FlowFile Record로 출력하는 데 사용하는 형식이에요.
Processing Strategy Kafka Records를 처리하고 직렬화된 출력을 FlowFiles로 작성하는 전략이에요.
Record Reader 수신 Kafka 메시지에 사용할 Record Reader예요.
Record Writer 출력 FlowFiles를 직렬화하는 데 사용할 Record Writer예요.
Separate By Key 이 속성을 활성화하면 두 Kafka 메시지의 키가 모두 동일한 경우에만 같은 FlowFile에 추가돼요.
Topic Format 제공된 Topics가 쉼표로 구분된 이름 목록인지 단일 정규 표현식인지 지정해요.
Topics 프로세서가 Kafka Records를 소비하는 Kafka Topics의 이름 또는 패턴이에요. 쉼표로 구분해 여러 개를 제공할 수 있어요.
auto.offset.reset 이전 소비자 오프셋이 없을 때 적용되는 자동 오프셋 구성으로, Kafka auto.offset.reset 속성에 해당해요.

관계 (Relationships)

이름 설명
success 하나 이상의 직렬화된 Kafka Records를 담은 FlowFiles예요.

기록하는 속성 (Writes attributes)

이름 설명
record.count 수신된 레코드 수예요.
mime.type 구성된 Record Writer가 제공하는 MIME Type이에요.
kafka.count 두 개 이상일 때 기록된 메시지 수예요.
kafka.key 메시지의 키가 있고 단일 메시지일 때의 키예요. 키가 인코딩되는 방식은 'Key Attribute Encoding' 속성 값에 따라 달라져요.
kafka.offset 토픽의 파티션에서 메시지의 오프셋이에요.
kafka.timestamp 토픽의 파티션에서 메시지의 타임스탬프예요.
kafka.partition 메시지 또는 메시지 묶음이 속한 토픽의 파티션이에요.
kafka.topic 메시지 또는 메시지 묶음이 속한 토픽이에요.
kafka.tombstone 소비된 메시지가 tombstone 메시지이면 true로 설정돼요.

참고 (See also)

  • com.snowflake.openflow.runtime.processors.kafka.PublishKafka

더 알아보기 (Learn more)