PublishKafka

PublishKafka

FlowFile 내용을 메시지 또는 개별 레코드로 Kafka Producer API를 사용해 Apache Kafka에 보내는 프로세서예요.

출처: Snowflake 문서 — PublishKafka

본문

번들 (Bundle)

com.snowflake.openflow.runtime | runtime-kafka-nar

설명

Kafka Producer API를 사용해 FlowFile 내용을 메시지 또는 개별 레코드로 Apache Kafka에 보내요. 보낼 메시지는 개별 FlowFile일 수도, 사용자 지정 구분자(예: 개행)로 구분될 수도, 또는 구성된 Record Reader로 읽을 수 있는 레코드 지향 데이터일 수도 있어요. 메시지를 가져오는 보완 NiFi 프로세서는 ConsumeKafka예요.

태그

apache, avro, csv, json, kafka, logs, message, openflow, pubsub, put, record, send

입력 요구 사항 (Input Requirement)

REQUIRED — 입력 FlowFile이 필요해요.

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

false — 지원하지 않아요.

속성 (Properties)

속성 설명
Failure Strategy 데이터를 Kafka에 게시할 수 없을 때 FlowFile을 처리하는 방식을 지정해요.
FlowFile Attribute Header Pattern 모든 FlowFile 속성 이름에 대해 일치시킬 정규식이에요. 이름이 패턴과 일치하는 모든 속성은 Kafka 메시지에 헤더로 추가돼요. 지정하지 않으면 FlowFile 속성이 헤더로 추가되지 않아요.
Header Encoding Kafka Record Header로 추가되는 속성에 대해 헤더 직렬화에 사용할 문자 인코딩을 나타내요.
Kafka Connection Service Kafka Records를 게시하기 위해 Kafka Broker에 대한 연결을 제공해요.
Kafka Key 메시지에 사용할 Key예요. 지정하지 않으면 FlowFile 속성 'kafka.key'가 존재할 때 메시지 키로 사용돼요. Kafka key를 설정하면서 동시에 demarcating 하면 같은 키를 가진 Kafka 메시지가 많아질 수 있어 주의하세요. 보통 Kafka는 메시지와 키의 고유성을 강제하거나 가정하지 않으므로 문제가 되지 않아요. 그래도 demarcator와 Kafka key를 동시에 설정하면 Kafka에서 데이터 손실 위험이 있어요. Kafka에서 토픽 압축 동안 메시지는 이 키를 기준으로 중복 제거돼요.
Kafka Key Attribute Encoding 방출된 FlowFile에는 'kafka.key'라는 속성이 있어요. 이 속성은 해당 속성 값을 인코딩하는 방법을 지정해요.
Message Demarcator 단일 FlowFile 안의 여러 메시지를 구분하는 데 사용할 문자열(UTF-8로 해석)을 지정해요. 지정하지 않으면 FlowFile 전체 내용이 단일 메시지로 사용돼요. 지정하면 FlowFile 내용이 이 구분자로 분할되고 각 구간이 별도의 Kafka 메시지로 전송돼요. 'new line' 같은 특수 문자를 입력하려면 OS에 따라 CTRL+Enter 또는 Shift+Enter를 사용하세요.
Message Key Field Kafka 메시지의 Key로 사용해야 할 입력 레코드의 필드 이름이에요.
Publish Strategy 들어오는 FlowFile 레코드를 Kafka에 게시하는 데 사용하는 형식이에요.
Record Key Writer 나가는 FlowFile에 사용할 Record Key Writer예요.
Record Metadata Strategy Record의 메타데이터(토픽 및 파티션)를 Record의 메타데이터 필드에서 가져올지, 구성된 Topic Name 및 Partition / Partitioner class 속성에서 가져올지 지정해요.
Record Reader 들어오는 FlowFile에 사용할 Record Reader예요.
Record Writer Kafka로 보내기 전에 데이터를 직렬화하는 데 사용할 Record Writer예요.
Topic Name 프로세서가 Kafka Records를 게시하는 Kafka Topic의 이름이에요.
Transactional ID Prefix KafkaProducer 구성 transactional.id가 생성된 UUID가 되고 구성된 문자열이 접두사로 붙는다는 것을 지정해요.
Transactions Enabled Kafka와 통신할 때 트랜잭션 보장을 제공할지 여부를 지정해요. Kafka로 데이터를 보내는 데 문제가 있고 이 속성이 false면 이미 Kafka에 전송된 메시지는 계속 진행되어 소비자에게 전달돼요. true면 Kafka 트랜잭션이 롤백되어 해당 메시지가 소비자에게 제공되지 않아요. true로 설정하려면 [Delivery Guarantee] 속성을 [Guarantee Replicated Delivery]로 설정해야 해요.
acks 메시지가 Kafka에 전송되었음을 보장하기 위한 요구 사항을 지정해요. Kafka Client의 acks 속성에 해당해요.
compression.type Kafka로 전송되는 레코드의 압축 전략을 지정해요. Kafka Client의 compression.type 속성에 해당해요.
max.request.size 요청의 최대 크기(바이트)예요. Kafka Client의 max.request.size 속성에 해당해요.
partition 레코드의 Kafka Partition 목적지를 지정해요.
partitioner.class 메시지의 파티션 ID를 계산하는 데 사용할 클래스를 지정해요. Kafka Client의 partitioner.class 속성에 해당해요.

관계 (Relationships)

이름 설명
failure Kafka에 보낼 수 없는 모든 FlowFile이 이 관계로 라우팅돼요.
success 모든 내용이 Kafka로 전송된 FlowFile이 이 관계로 이동해요.

쓰기 속성 (Writes attributes)

이름 설명
msg.count 이 FlowFile에 대해 Kafka로 전송된 메시지 수예요. 이 속성은 success로 라우팅되는 FlowFile에만 추가돼요.

더 보기 (See also)

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

더 알아보기 (Learn more)