ConsumeAMQP 프로세서
ConsumeAMQP 프로세서
이 문서는 Snowflake OpenFlow의 ConsumeAMQP 프로세서에 대한 참조 문서예요. AMQP 0.9.1 프로토콜을 사용해 AMQP 브로커에서 메시지를 소비하고, 각 메시지를 별도의 FlowFile로 내보내요.
출처: Snowflake 문서
본문
기능 — 일반 제공 (Generally Available)
Openflow Snowflake 배포는 AWS, Azure, GCP Commercial 리전의 모든 계정에서 사용할 수 있어요.
Openflow BYOC 배포는 AWS Commercial 리전의 모든 계정에서 사용할 수 있어요.
번들 (Bundle)
| 그룹 | NAR |
|---|---|
| org.apache.nifi | nifi-amqp-nar |
설명 (Description)
AMQP 0.9.1 프로토콜을 사용해 AMQP 브로커에서 AMQP 메시지를 소비해요. AMQP 브로커로부터 수신된 각 메시지는 자체 FlowFile로 발행되어 'success' 관계로 전달돼요.
태그 (Tags)
amqp, consume, get, message, rabbit, receive
입력 요구사항 (Input Requirement)
FORBIDDEN
민감한 동적 속성 지원 (Supports Sensitive Dynamic Properties)
아니요 (false)
속성 (Properties)
| 속성 | 설명 |
|---|---|
| AMQP Version | AMQP 버전이에요. 현재는 AMQP v0.9.1만 지원해요. |
| Auto-Acknowledge Messages | false(비자동 승인)이면 FlowFiles를 success로 전달하고 NiFi 세션을 커밋한 후 프로세서가 메시지를 승인해요. 비자동 승인 모드는 'at-least-once'(최소 한 번) 전달 의미론을 제공해요. true(자동 승인)이면 AMQP 클라이언트에 전달된 메시지가 전송 직후 AMQP 브로커에 의해 자동 승인돼요. 일반적으로 처리량은 더 좋지만 AMQP 브로커, NiFi 또는 프로세서가 재시작되거나 충돌하면 메시지가 손실될 수 있어요. 자동 승인 모드는 'at-most-once'(최대 한 번) 전달 의미론을 제공하며, 메시지 손실이 허용되는 경우에만 권장돼요. |
| Batch Size | 단일 세션에서 처리할 최대 메시지 수예요. 이만큼의 메시지를 수신하거나(또는 더 이상 메시지를 사용할 수 없게 되면) 수신된 메시지가 'success' 관계로 전달되고 메시지가 AMQP 브로커에 승인돼요. 이 값을 크게 설정하면 특히 매우 작은 메시지에서 성능이 좋아질 수 있지만, NiFi가 갑작스럽게 재시작될 때 메시지가 더 많이 중복될 수 있어요. |
| Brokers | 알려진 AMQP 브로커의 쉼표로 구분된 목록으로, 형식은 호스트:포트예요(예: localhost:5672). 이 값이 설정되면 Host Name과 Port는 무시돼요. 같은 AMQP 클러스터의 호스트만 포함해야 해요. |
| Client Certificate Authentication Enabled | 사용자 이름/비밀번호 대신 SSL 인증서를 사용해 인증해요. |
| Header Key Prefix | FlowFile 속성으로 추가될 헤더 키 앞에 붙일 텍스트예요. 프로세서는 이 속성 값에 '.'을 추가해요. |
| Header Output Format | 수신 메시지의 헤더를 출력하는 방식을 정의해요. |
| Header Separator | String 형태의 헤더에서 키-값을 구분하는 데 사용하는 문자예요. 값은 반드시 한 글자여야 해요. |
| Host Name | AMQP 브로커의 네트워크 주소예요(예: localhost). Brokers가 설정되면 이 속성은 무시돼요. |
| Max Inbound Message Body Size | 인바운드(수신) 메시지의 최대 본문 크기예요. |
| Password | 인증 및 권한 부여에 사용되는 비밀번호예요. |
| Port | AMQP 브로커의 포트를 식별하는 숫자 값이에요(예: 5671). Brokers가 설정되면 이 속성은 무시돼요. |
| Prefetch Count | 소비자에 대한 미승인(unacknowledged) 메시지의 최대 수예요. 소비자가 이만큼의 미승인 메시지를 가지면, 이미 전달된 메시지 중 일부를 승인할 때까지 AMQP 브로커는 새 메시지를 보내지 않아요. 허용 값: 0~65535. 0은 제한 없음을 의미해요. |
| Queue | 메시지를 소비할 기존 AMQP Queue의 이름이에요. 보통 AMQP 관리자가 미리 정의해요. |
| Remove Curly Braces | true이면 헤더의 중괄호(Curly Braces)가 자동으로 제거돼요. |
| SSL Context Service | TLS/SSL 연결에 클라이언트 인증서 정보를 제공하는 데 사용할 SSL Context Service예요. |
| Username | 인증 및 권한 부여에 사용되는 사용자 이름이에요. |
| Virtual Host | 보안 강화를 위해 AMQP 시스템을 분리하는 Virtual Host 이름이에요. |
관계 (Relationships)
| 이름 | 설명 |
|---|---|
| success | AMQP 큐에서 수신된 모든 FlowFiles가 이 관계로 라우팅돼요. |
기록하는 속성 (Writes attributes)
| 이름 | 설명 |
|---|---|
| amqp$appId | AMQP 메시지의 App ID 필드예요. |
| amqp$contentEncoding | AMQP 메시지가 보고한 Content Encoding이에요. |
| amqp$contentType | AMQP 메시지가 보고한 Content Type이에요. |
| amqp$headers | AMQP 메시지에 있는 헤더예요. 프로세서가 이 속성을 출력하도록 구성된 경우에만 추가돼요. |
| . | 각 메시지 헤더는 프로세서가 헤더를 속성으로 출력하도록 구성된 경우 이 속성 이름으로 삽입돼요. |
| amqp$deliveryMode | 메시지의 Delivery Mode에 대한 숫자 표시예요. |
| amqp$priority | 메시지 우선순위예요. |
| amqp$correlationId | 메시지의 Correlation ID예요. |
| amqp$replyTo | 메시지의 Reply-To 필드 값이에요. |
| amqp$expiration | 메시지 만료(Expiration)예요. |
| amqp$messageId | 메시지의 고유 ID예요. |
| amqp$timestamp | epoch 이후 밀리초 수로 표시된 메시지 타임스탬프예요. |
| amqp$type | 메시지 유형이에요. |
| amqp$userId | 사용자의 ID예요. |
| amqp$clusterId | AMQP 클러스터의 ID예요. |
| amqp$routingKey | AMQP 메시지의 routingKey예요. |
| amqp$exchange | AMQP 메시지를 받은 exchange예요. |