ConsumeMQTT 프로세서
ConsumeMQTT 프로세서
이 문서는 Snowflake OpenFlow의 ConsumeMQTT 프로세서에 대한 참조 문서예요. 토픽에 구독하고 MQTT 브로커로부터 메시지를 수신해요.
출처: Snowflake 문서
본문
기능 — 일반 제공 (Generally Available)
Openflow Snowflake 배포는 AWS, Azure, GCP Commercial 리전의 모든 계정에서 사용할 수 있어요.
Openflow BYOC 배포는 AWS Commercial 리전의 모든 계정에서 사용할 수 있어요.
번들 (Bundle)
| 그룹 | NAR |
|---|---|
| org.apache.nifi | nifi-mqtt-nar |
설명 (Description)
토픽에 구독하고 MQTT 브로커로부터 메시지를 수신해요.
태그 (Tags)
IOT, MQTT, consume, listen, subscribe
입력 요구사항 (Input Requirement)
FORBIDDEN
민감한 동적 속성 지원 (Supports Sensitive Dynamic Properties)
아니요 (false)
속성 (Properties)
| 속성 | 설명 |
|---|---|
| Broker URI | MQTT 브로커에 연결하는 데 사용할 URI예요(예: tcp://localhost:1883). 'tcp', 'ssl', 'ws', 'wss' 스킴이 지원돼요. 'ssl'을 사용하려면 SSL Context Service 속성을 설정해야 해요. 쉼표로 구분된 URI 목록이 설정되면(예: tcp://localhost:1883,tcp://localhost:1884) 프로세서는 연결 실패 시 라운드 로빈 알고리즘을 사용해 브로커에 연결해요. |
| Client ID | 사용할 MQTT 클라이언트 ID예요. 설정하지 않으면 UUID가 생성돼요. |
| Connection Timeout (seconds) | MQTT 서버에 대한 네트워크 연결이 수립되기까지 클라이언트가 기다리는 최대 시간 간격이에요. 기본 시간 제한은 30초예요. 값 0은 시간 제한 처리를 비활성화하며, 네트워크 연결이 성공하거나 실패할 때까지 클라이언트가 기다린다는 뜻이에요. |
| Group ID | 사용할 MQTT consumer group ID예요. group ID가 설정되지 않으면 클라이언트가 개별 소비자로 연결돼요. |
| Keep Alive Interval (seconds) | 보내거나 받는 메시지 사이의 최대 시간 간격을 정의해요. 이를 통해 TCP/IP 시간 제한을 기다리지 않고도 서버를 더 이상 사용할 수 없는지 클라이언트가 감지할 수 있어요. 클라이언트는 각 keep alive 주기 내에 최소 하나의 메시지가 네트워크를 통해 전송되도록 보장해요. 해당 기간 동안 데이터 관련 메시지가 없으면 클라이언트는 매우 작은 "ping" 메시지를 보내고 서버가 이를 승인해요. 값 0은 클라이언트의 keepalive 처리를 비활성화해요. |
| Last Will Message | 클라이언트의 Last Will로 보낼 메시지예요. |
| Last Will QoS Level | Last Will Message를 게시할 때 사용할 QoS 레벨이에요. |
| Last Will Retain | 클라이언트의 Last Will를 보존할지 여부예요. |
| Last Will Topic | 클라이언트의 Last Will를 보낼 토픽이에요. |
| MQTT Specification Version | 브로커에 연결할 때의 MQTT 사양 버전이에요. 자세한 내용은 허용 값 설명을 참고해요. |
| Max Queue Size | MQTT 메시지는 프로세서가 실행되도록 스케줄링된 빈도와 관계없이 항상 토픽의 구독자에게 전송돼요. 'Run Schedule'이 이 프로세서로 메시지가 도착하는 속도보다 크게 뒤처지면 이 프로세서의 내부 큐에 백업이 발생할 수 있어요. 이 속성은 이 프로세서가 한 번에 내부 큐에 메모리에 보관할 최대 메시지 수를 지정해요. 이 데이터는 NiFi 재시작 시 손실될 수 있어요. |
| Password | 브로커에 연결할 때 사용할 비밀번호예요. |
| Quality of Service(QoS) | 메시지를 수신할 QoS(Quality of Service)예요. '0', '1', '2' 값을 허용하며, '0'은 'at most once', '1'은 'at least once', '2'는 'exactly once'예요. |
| SSL Context Service | TLS/SSL 연결에 클라이언트 인증서 정보를 제공하는 데 사용할 SSL Context Service예요. |
| Session Expiry Interval | 이 간격 후에 브로커가 클라이언트를 만료시키고 세션 상태를 정리해요. |
| Session state | 새로 시작할지 이전 흐름을 재개할지 여부예요. 자세한 내용은 허용 값 설명을 참고해요. |
| Topic Filter | 구독할 토픽을 지정하는 MQTT 토픽 필터예요. |
| Username | 브로커에 연결할 때 사용할 사용자 이름이에요. |
| add-attributes-as-fields | 이 속성을 true로 설정하면 각 레코드에 기본 필드 _topic, _qos, _isDuplicate, _isRetained가 추가돼요. |
| message-demarcator | 이 속성을 사용하면 여러 메시지를 담은 FlowFiles를 출력하는 옵션을 사용할 수 있어요. 여러 메시지를 구분하는 데 사용할 문자열(UTF-8로 해석)을 제공할 수 있어요. 선택 속성이며, 제공하지 않고 Record Reader/Writer를 정의하지 않으면 수신된 각 메시지가 단일 FlowFile이 돼요. 'new line' 같은 특수 문자를 입력하려면 OS에 따라 CTRL+Enter 또는 Shift+Enter를 사용해요. |
| record-reader | 수신된 MQTT Messages를 Records로 파싱하는 데 사용할 Record Reader예요. |
| record-writer | FlowFile에 쓰기 전에 Records를 직렬화하는 데 사용할 Record Writer예요. |
관계 (Relationships)
| 이름 | 설명 |
|---|---|
| Message | MQTT 메시지 출력이에요. |
| parse.failure | 구성된 Record Reader로 메시지를 파싱할 수 없으면 메시지 내용이 자체 FlowFile로 이 Relationship에 라우팅돼요. |
기록하는 속성 (Writes attributes)
| 이름 | 설명 |
|---|---|
| record.count | 수신된 레코드 수예요. |
| mqtt.broker | 메시지 소스였던 MQTT 브로커예요. |
| mqtt.topic | 메시지가 수신된 MQTT 토픽이에요. |
| mqtt.qos | 이 메시지의 서비스 품질(QoS)이에요. |
| mqtt.isDuplicate | 이 메시지가 이미 수신한 메시지의 중복일 가능성이 있는지 여부예요. |
| mqtt.isRetained | 이 메시지가 현재 게시자로부터 온 것인지, 아니면 서버가 해당 토픽에 게시된 마지막 메시지로 "retained"(보존)한 것인지 여부예요. |
참고 (See also)
org.apache.nifi.processors.mqtt.PublishMQTT