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

더 알아보기 (Learn more)