PublishKafka
PublishKafka
FlowFile 내용을 메시지 또는 개별 레코드로 Kafka Producer API를 사용해 Apache Kafka에 보내는 프로세서예요.
본문
번들 (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