Pulsar Functions 개념
Pulsar Functions 개념 (Pulsar Functions Concepts)
Pulsar Functions는 입력 토픽의 메시지를 처리해 출력 토픽으로 보내는 경량 컴퓨팅 모델이에요. 함수 하나가 어떻게 실행되는지, 전달 보장(processing guarantee)이 무엇인지 이해하면 효과적으로 만들고 배포할 수 있어요. 이 글에서는 함수의 핵심 개념(FQFN, 함수 인스턴스, function worker, 런타임)과 전달 시맨틱, 컨텍스트 객체를 정리해 드릴게요.
출처: 문서
본문
완전 정규 함수 이름 (FQFN)
각 함수는 지정된 테넌트, 네임스페이스, 함수 이름을 가진 완전 정규 함수 이름(FQFN, Fully Qualified Function Name)을 가져요. FQFN 덕분에 서로 다른 네임스페이스에 같은 함수 이름으로 여러 함수를 만들 수 있어요.
FQFN은 다음과 같아요:
tenant/namespace/name
함수 인스턴스
함수 인스턴스는 함수 실행 프레임워크의 핵심 요소이며 다음 요소로 구성돼요:
- 서로 다른 입력 토픽에서 메시지를 소비하는 소비자 모음.
- 함수를 호출하는 실행기(executor).
- 함수의 결과를 출력 토픽으로 보내는 프로듀서.
다음 그림은 함수 인스턴스의 내부 워크플로를 보여줘요.
함수는 여러 인스턴스를 가질 수 있고, 각 인스턴스는 함수의 복사본 하나를 실행해요. 구성 파일에서 인스턴스 수를 지정할 수 있어요.
함수 인스턴스 안의 소비자들은 구독 유형을 기반으로 여러 인스턴스 간 부하 분산을 활성화하기 위해 FQFN을 구독자 이름으로 사용해요. 구독 유형은 함수 수준에서 지정할 수 있어요.
각 함수는 FQFN으로 별도의 상태 저장소를 가져요. 중간 결과를 BookKeeper에 영구 저장하기 위해 상태 인터페이스를 지정할 수 있어요. 다른 사용자가 함수의 상태를 조회하고 이 결과를 추출할 수 있어요.
Function worker
Function worker는 Pulsar Functions의 클러스터 모드 배포에서 개별 함수를 모니터링, 오케스트레이션, 실행하는 논리 구성 요소예요.
function worker 안에서 각 함수 인스턴스는 선택한 구성에 따라 스레드나 프로세스로 실행될 수 있어요. 또는 Kubernetes 클러스터를 사용할 수 있다면 함수를 Kubernetes 안의 StatefulSet으로 스폰할 수 있어요. 자세한 내용은 Set up function workers를 보세요.
다음 그림은 function worker의 내부 아키텍처와 워크플로를 보여줘요.
Function worker는 워커 노드 클러스터를 형성하고 워크플로는 다음과 같아요.
- 사용자가 REST 서버에 요청을 보내 함수 인스턴스를 실행해요.
- REST 서버가 요청에 응답하고 요청을 함수 메타데이터 관리자에게 전달해요.
- 함수 메타데이터 관리자가 요청 갱신 내용을 함수 메타데이터 토픽에 써요. 또한 모든 메타데이터 관련 메시지를 추적하고 함수 메타데이터 토픽을 사용해 함수의 상태 갱신을 영구 저장해요.
- 함수 메타데이터 관리자가 함수 메타데이터 토픽에서 갱신 내용을 읽고, 스케줄 관리자에게 할당(assignment)을 계산하도록 트리거해요.
- 스케줄 관리자가 할당 갱신 내용을 할당 토픽에 써요.
- 함수 런타임 관리자가 할당 토픽을 수신하고 할당 갱신을 읽어, 모든 워커에 대한 모든 할당의 전역 뷰를 포함하는 내부 상태를 갱신해요. 갱신이 특정 워커의 할당을 변경하면 함수 런타임 관리자는 함수 인스턴스의 실행을 시작하거나 중지해 새 할당을 구체화해요.
- 멤버십 관리자가 조정 토픽에 요청해 리드 워커를 선출해요. 모든 워커는 failover 구독으로 조정 토픽에 구독하지만, 활성 워커가 리더가 되어 할당을 수행하므로 이 토픽에 활성 소비자가 하나만 있게 보장돼요.
- 멤버십 관리자가 조정 토픽에서 갱신 내용을 읽어요.
함수 런타임
함수 인스턴스는 런타임 안에서 호출되고, 여러 인스턴스가 병렬로 실행될 수 있어요. Pulsar는 배포 유연성을 최대화하기 위해 비용과 격리 보장이 다른 세 가지 함수 런타임 유형을 지원해요. 필요에 따라 이 중 하나를 사용해 함수를 실행할 수 있어요. 자세한 내용은 함수 런타임 구성을 보세요.
다음 표는 세 가지 함수 런타임 유형을 설명해요.
| 유형 | 설명 |
|---|---|
| 스레드 런타임 | 각 인스턴스가 스레드로 실행돼요. 스레드 모드의 코드는 Java로 작성되므로 Java 인스턴스에만 적용돼요. 함수가 스레드 모드로 실행되면 function worker와 같은 Java 가상 머신(JVM)에서 실행돼요. |
| 프로세스 런타임 | 각 인스턴스가 프로세스로 실행돼요. 함수가 프로세스 모드로 실행되면 function worker가 실행되는 것과 같은 머신에서 실행돼요. |
| Kubernetes 런타임 | 함수가 워커에 의해 Kubernetes StatefulSet으로 제출되고 각 함수 인스턴스가 pod로 실행돼요. Pulsar는 함수를 실행할 때 Kubernetes StatefulSet과 서비스에 레이블을 추가하도록 지원해, 대상 Kubernetes 객체를 선택하기 쉽게 해줘요. |
처리 보장과 구독 유형
Pulsar는 함수에 적용할 수 있는 세 가지 서로 다른 메시지 전달 시맨틱을 제공해요. 서로 다른 전달 시맨틱 구현은 **확인 시점 노드(ack time node)**에 따라 결정돼요.
| 전달 시맨틱 | 설명 | 채택된 구독 유형 |
|---|---|---|
| At-most-once 전달 | 함수에 보내진 각 메시지는 최선으로 처리돼요. 메시지가 처리될지 여부는 보장되지 않아요. 이 시맨틱을 선택하면 autoAck 구성이 true로 설정되어야 해요. 그렇지 않으면 시작이 실패해요 (autoAck 구성은 향후 릴리스에서 폐기될 예정이에요). 확인 시점 노드: 함수 처리 전. |
Shared |
| At-least-once 전달 (기본값) | 함수에 보내진 각 메시지는 (처리 실패나 재전달의 경우) 두 번 이상 처리될 수 있어요. --processing-guarantees 플래그를 지정하지 않고 함수를 만들면 함수는 at-least-once 전달 보장을 제공해요. 확인 시점 노드: 출력으로 메시지 보낸 후. |
Shared |
| Effectively-once 전달 | 함수에 보내진 각 메시지는 두 번 이상 처리될 수 있지만 출력은 하나만 있어요. 중복 메시지는 무시돼요. Effectively once는 at-least-once 처리 위에서 달성되며 서버 측 중복 제거가 보장돼요. 즉 상태 갱신이 두 번 일어날 수 있지만, 같은 상태 갱신은 한 번만 적용되고 다른 중복 상태 갱신은 서버 측에서 버려져요. 확인 시점 노드: 출력으로 메시지 보낸 후. |
Failover |
| Manual 전달 | 이 시맨틱을 선택하면 프레임워크가 어떤 확인 작업도 수행하지 않아요. 함수 안에서 context.getCurrentRecord().ack() 메서드를 호출해 확인 작업을 직접 수행해야 해요. 확인 시점 노드: 함수 메서드 안에서 사용자 정의. |
Shared |
팁
- 기본적으로 Pulsar Functions는 at-least-once 전달 보장을 제공해요.
--processingGuarantees플래그에 값을 제공하지 않고 함수를 만들면 함수가 at-least-once 보장을 제공해요.- Exclusive 구독 유형은 Pulsar Functions에서 사용할 수 없어요. 그 이유는:
- 인스턴스가 하나뿐이면 exclusive는 failover와 같아요.
- 인스턴스가 여러 개면 함수를 재시작할 때 exclusive가 충돌·재시작할 수 있어요. 이 경우 exclusive는 failover와 같지 않아요. 마스터 소비자가 연결을 끊으면 모든 미확인 및 이후 메시지가 대기열의 다음 소비자에게 전달되기 때문이에요.
- 구독 유형을 shared에서 key_shared로 변경하려면
pulsar-admin에서—retain-key-ordering옵션을 사용할 수 있어요.
함수를 만들 때 함수의 처리 보장을 설정할 수 있어요. 다음 명령은 effectively-once 보장이 적용된 함수를 만들어요.
bin/pulsar-admin functions create \
--name my-effectively-once-function \
--processing-guarantees EFFECTIVELY_ONCE \
# Other function configs
update 명령으로 함수에 적용된 처리 보장을 변경할 수 있어요.
bin/pulsar-admin functions update \
--processing-guarantees ATMOST_ONCE \
# Other function configs
컨텍스트 (Context)
Java, Python, Go SDK는 함수가 사용할 수 있는 컨텍스트 객체에 대한 접근을 제공해요. 이 컨텍스트 객체는 함수에 다양한 정보와 기능을 제공해요:
- 함수의 이름과 ID.
- 메시지의 메시지 ID. 각 메시지에는 자동으로 ID가 할당돼요.
- 메시지의 키, 이벤트 시간, 속성, 파티션 키.
- 메시지가 전송된 토픽의 이름.
- 함수와 연관된 모든 입력 토픽 및 출력 토픽의 이름.
- SerDe에 사용된 클래스의 이름.
- 함수와 연관된 테넌트와 네임스페이스.
- 함수를 실행하는 함수 인스턴스의 ID.
- 함수의 버전.
- 함수가 사용하는 로거 객체. 로그 메시지를 만드는 데 사용돼요.
- CLI를 통해 제공된 임의 사용자 구성 값에 대한 접근.
- 지표를 기록하는 인터페이스.
- 상태 저장소에서 상태를 저장하고 검색하는 인터페이스.
- 임의의 토픽에 새 메시지를 게시하는 함수.
- 처리 중인 메시지를 확인하는 함수 (자동 확인이 비활성화된 경우).
- (Java) Pulsar admin 클라이언트를 가져오는 함수.
- (Java) Context와 입력 Record의 기본값으로 Return할 Record를 만드는 함수.
함수 메시지 유형
Pulsar Functions는 입력으로 바이트 배열을 받고 출력으로 바이트 배열을 내뱉어요. 다음 두 가지 방식 중 하나로 타입 있는 함수를 작성하고 메시지를 타입에 바인딩할 수 있어요:
- 스키마 레지스트리 (Schema Registry)
- SerDe
윈도우 함수 (Window function)
참고 현재 윈도우 함수는 Java에서만 사용할 수 있으며,
MANUAL및Effectively-once전달 시맨틱을 지원하지 않아요.
윈도우 함수는 데이터 윈도우, 즉 이벤트 스트림의 유한 부분집합에 걸쳐 계산을 수행하는 함수예요. 아래에 설명된 대로 스트림은 함수를 적용할 수 있는 "버킷"으로 나뉘어요.
함수의 데이터 윈도우 정의는 두 가지 정책을 포함해요:
- 제거 정책(Eviction policy): 윈도우에 수집되는 데이터 양을 제어해요.
- 트리거 정책(Trigger policy): 제거 정책에 따라 윈도우에 수집된 모든 데이터를 처리하기 위해 함수가 언제 트리거되고 실행되는지 제어해요.
트리거 정책과 제거 정책 모두 시간 또는 개수에 의해 구동돼요.
팁 처리 시간(processing time)과 이벤트 시간(event time)이 모두 지원돼요.
- 처리 시간은 함수 인스턴스가 윈도우를 만들고 처리하는 벽(wall) 시간을 기준으로 정의돼요. 윈도우 완전성 판단이 간단하며 데이터 도착 순서 불일치를 걱정할 필요가 없어요.
- 이벤트 시간은 이벤트 레코드와 함께 오는 타임스탬프를 기준으로 정의돼요. 이벤트 시간의 정확성은 보장하지만 더 많은 데이터 버퍼링과 제한된 완전성 보장을 제공해요.
윈도우 유형
인접한 두 윈도우가 공통 이벤트를 공유할 수 있는지 여부에 따라 윈도우는 다음 두 가지 유형으로 나뉘어요:
- Tumbling window
- Sliding window
Tumbling window
Tumbling window는 지정된 시간 길이 또는 개수의 윈도우에 요소를 할당해요. tumbling window의 제거 정책은 항상 윈도우가 가득 찬 것을 기준으로 해요. 따라서 개수 기반 또는 시간 기반 중 하나의 트리거 정책만 지정하면 돼요.
개수 기반 트리거 정책의 tumbling window에서는 아래 예시처럼 트리거 정책이 2로 설정되어 있어요. 시간과 무관하게 윈도우에 두 항목이 있으면 각 함수가 트리거되고 실행돼요.
반대로 아래 예시처럼 tumbling window의 윈도우 길이가 10초라면, 윈도우에 이벤트가 몇 개든 10초 시간 간격이 경과했을 때 함수가 트리거되는 것을 의미해요.
Sliding window
Sliding window 방식은 제거 정책을 설정해 처리에 보존할 데이터 양을 제한하고, 슬라이딩 간격으로 트리거 정책을 설정해 고정 윈도우 길이를 정의해요. 슬라이딩 간격이 윈도우 길이보다 작으면 데이터가 겹치는데, 이는 인접한 윈도우에 동시에 속하는 데이터가 계산에 두 번 이상 사용된다는 뜻이에요.
아래 예시처럼 윈도우 길이가 2초라면, 2초보다 오래된 데이터는 제거되고 계산에 사용되지 않아요. 슬라이딩 간격이 1초로 구성되어 있으므로, 함수가 매초 실행되어 전체 윈도우 길이 안의 데이터를 처리해요.
더 알아보기 (Learn more)
- Functions 개요 — Pulsar Functions가 무엇인지 이해해요.
- Functions CLI — 함수를 만드는 CLI 명령을 익혀요.
- Functions 배포 — 클러스터 모드로 함수를 배포해요.
- 메시징 개념 — 구독 유형을 익혀요.