ConsumeAzureEventHub 프로세서

ConsumeAzureEventHub 프로세서

이 문서는 Snowflake OpenFlow의 ConsumeAzureEventHub 프로세서에 대한 참조 문서예요. 마이크로소프트 Azure Event Hubs에서 메시지를 수신하며, 체크포인트(checkpoint)를 이용해 일관된 이벤트 처리를 보장해요.

출처: Snowflake 문서

본문

기능 — 일반 제공 (Generally Available)

Openflow Snowflake 배포는 AWS, Azure, GCP Commercial 리전의 모든 계정에서 사용할 수 있어요.

Openflow BYOC 배포는 AWS Commercial 리전의 모든 계정에서 사용할 수 있어요.

번들 (Bundle)

그룹 NAR
org.apache.nifi nifi-azure-nar

설명 (Description)

일관된 이벤트 처리를 보장하기 위해 체크포인트를 사용해 마이크로소프트 Azure Event Hubs에서 메시지를 수신해요. 체크포인트 추적은 메시지를 여러 번 소비하는 것을 방지하고, 일시적인 네트워크 오류 시 처리를 안정적으로 재개할 수 있게 해줘요. 체크포인트 추적에는 외부 저장소가 필요하며, Azure Event Hubs에서 메시지를 소비하는 권장 방식이에요. 클러스터 환경에서는 ConsumeAzureEventHub 프로세서 인스턴스들이 하나의 consumer group을 형성하고 메시지가 클러스터 노드들에 분산돼요(각 메시지는 하나의 클러스터 노드에서만 처리돼요).

태그 (Tags)

azure, cloud, eventhub, events, microsoft, streaming, streams

입력 요구사항 (Input Requirement)

FORBIDDEN

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

아니요 (false)

속성 (Properties)

속성 설명
Batch Size NiFi 세션 내에서 처리할 메시지 수예요. 이 매개변수는 처리량과 일관성에 영향을 줘요. NiFi는 이 수만큼의 메시지를 처리한 후 세션과 Event Hubs 체크포인트를 커밋해요. NiFi 세션이 커밋됐는데 Event Hubs 체크포인트 생성에 실패하면 같은 메시지를 다시 받을 수 있어요. 숫자가 높을수록 처리량은 높아지지만 일관성은 떨어질 수 있어요.
Checkpoint Strategy 각 파티션에 대한 파티션 소유권과 체크포인트 정보를 저장하고 검색하는 데 사용할 전략을 지정해요.
Consumer Group 사용할 consumer group의 이름이에요.
Event Hub Name 메시지를 가져올 event hub의 이름이에요.
Event Hub Namespace Azure Event Hubs가 할당된 네임스페이스예요. 일반적으로 -ns와 같아요.
Initial Offset 체크포인트 저장소에 오프셋이 아직 저장되지 않았을 때 메시지 수신을 시작할 위치를 지정해요.
Message Receive Timeout 이 소비자가 Batch Size만큼 수신하기 위해 기다려야 하는 시간이에요.
Prefetch Count (설명 없음)
Record Reader 수신 메시지를 읽는 데 사용할 Record Reader예요. 스키마에 접근하려면 Expression Language '${eventhub.name}'으로 event hub 이름을 참조할 수 있어요.
Record Writer Records를 출력 FlowFile로 직렬화하는 데 사용할 Record Writer예요. 스키마에 접근하려면 Expression Language '${eventhub.name}'으로 event hub 이름을 참조할 수 있어요. 지정하지 않으면 각 메시지가 하나의 FlowFile을 만들어요.
Service Bus Endpoint 기본 windows.net 도메인에 없는 네임스페이스를 지원하기 위한 항목이에요.
Shared Access Policy Key 공유 액세스 정책의 키예요. 기본 키 또는 보조 키를 사용할 수 있어요.
Shared Access Policy Name 공유 액세스 정책의 이름이에요. 이 정책은 Listen 권한을 가져야 해요.
Storage Account Key event hub consumer group 상태를 저장할 Azure Storage 계정 키예요.
Storage Account Name event hub consumer group 상태를 저장할 Azure Storage 계정의 이름이에요.
Storage Container Name event hub consumer group 상태를 저장할 Azure Storage 컨테이너의 이름이에요. 지정하지 않으면 event hub 이름이 사용돼요.
Storage SAS Token Event Hub consumer group 상태를 저장할 Azure Storage SAS 토큰이에요. 항상 ? 문자로 시작해요.
Transport Type Azure Event Hubs와의 통신을 위한 Advanced Message Queuing Protocol Transport Type이에요.
Use Azure Managed Identity Azure VM/VMSS의 관리 ID(managed identity)를 사용할지 여부를 선택해요.
proxy-configuration-service 네트워크 요청을 프록시할 Proxy Configuration Controller Service를 지정해요.

상태 관리 (State management)

범위 설명
LOCAL Local state는 client id를 저장하는 데 사용돼요. 구성 요소 상태가 체크포인트 전략으로 구성되면 cluster state가 파티션 소유권과 체크포인트 정보를 저장하는 데 사용돼요.
CLUSTER Local state는 client id를 저장하는 데 사용돼요. 구성 요소 상태가 체크포인트 전략으로 구성되면 cluster state가 파티션 소유권과 체크포인트 정보를 저장하는 데 사용돼요.

관계 (Relationships)

이름 설명
success Event Hub에서 수신된 FlowFiles예요.

기록하는 속성 (Writes attributes)

이름 설명
eventhub.enqueued.timestamp 메시지가 event hub에 인큐된 시각(epoch 이후 UTC 밀리초)이에요.
eventhub.offset 메시지가 저장된 파티션 내 오프셋이에요.
eventhub.sequence 메시지와 연결된 시퀀스 번호예요.
eventhub.name 메시지를 가져온 event hub의 이름이에요.
eventhub.partition 메시지를 가져온 파티션의 이름이에요.
eventhub.property.* 이 메시지의 애플리케이션 속성이에요. 예: 'application'이면 'eventhub.property.application'이 돼요.

더 알아보기 (Learn more)