Pulsar Functions CLI 및 YAML 구성

Pulsar Functions CLI 및 YAML 구성

Pulsar Functions는 CLI로 만들고 관리하며, YAML 파일로 함수 구성을 정의할 수도 있어요. 이 글에서는 pulsar-admin CLI로 함수를 다루는 방법과, YAML 구성 파일의 모든 필드(명령 인자와 매핑 포함)를 표로 정리해 드릴게요.

출처: 문서

본문

Pulsar Functions용 Pulsar admin CLI

Pulsar admin 인터페이스는 CLI를 통해 Pulsar Functions를 만들고 관리할 수 있게 해줘요. 명령, 플래그, 설명을 포함한 최신·완전한 정보는 Pulsar admin CLI를 참조하세요.

Pulsar Functions용 YAML 구성

사전 정의된 YAML 파일로 함수를 구성할 수 있어요. 다음 표는 필수 필드와 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
runtimeFlags String N/A 런타임에 전달하려는 모든 플래그 (process 및 Kubernetes 런타임 전용).
tenant String --tenant 함수의 테넌트.
namespace String --namespace 함수의 네임스페이스.
name String --name 함수의 이름.
className String --classname 함수의 클래스 이름.
functionType String --function-type 내장 함수 유형.
inputs List<String> -i, --inputs 함수의 입력 토픽. 여러 토픽은 쉼표로 구분된 목록으로 지정할 수 있어요.
customSerdeInputs Map<String,String> --custom-serde-inputs 입력 토픽에서 SerDe 클래스 이름으로의 매핑.
topicsPattern String --topics-pattern 네임스페이스 아래 토픽 목록에서 소비하기 위한 토픽 패턴. 참고: --input--topic-pattern은 상호 배타적이에요. Java 함수의 경우 --custom-serde-inputs에 패턴에 대한 SerDe 클래스 이름을 추가해야 해요.
customSchemaInputs Map<String,String> --custom-schema-inputs 입력 토픽에서 스키마 속성으로의 매핑.
customSchemaOutputs Map<String,String> --custom-schema-outputs 출력 토픽에서 스키마 속성으로의 매핑.
inputSpecs Map<String,ConsumerConfig> --input-specs 입력에서 커스텀 구성으로의 매핑.
output String -o, --output 함수의 출력 토픽. 지정하지 않으면 출력이 기록되지 않아요.
producerConfig ProducerConfig --producer-config 프로듀서의 커스텀 구성.
outputSchemaType String -st, --schema-type 메시지 출력에 사용되는 내장 스키마 유형 또는 커스텀 스키마 클래스 이름.
outputSerdeClassName String --output-serde-classname 메시지 출력에 사용되는 SerDe 클래스.
logTopic String --log-topic 함수의 로그가 생성되는 토픽.
processingGuarantees String --processing-guarantees 함수에 적용되는 처리 보장(전달 시맨틱). 사용 가능한 값: ATLEAST_ONCE, ATMOST_ONCE, EFFECTIVELY_ONCE, MANUAL.
retainOrdering Boolean --retain-ordering 함수가 메시지를 순서대로 소비·처리하는지 여부.
retainKeyOrdering Boolean --retain-key-ordering 함수가 메시지를 키 순서대로 소비·처리하는지 여부.
batchBuilder String --batch-builder producerConfig.batchBuilder를 사용하세요. 참고: batchBuilder는 곧 코드에서 폐기될 예정이에요.
forwardSourceMessageProperty Boolean --forward-source-message-property 처리 중 입력 메시지의 속성을 출력 토픽으로 전달할지 여부. 값이 false로 설정되면 전달이 비활성화돼요.
userConfig Map<String,Object> --user-config 사용자 정의 구성 키/값.
secrets Map<String,Object> --secrets underlying secrets provider가 secret을 가져오는 방법을 캡슐화하는 객체로의 secretName 매핑.
runtime String N/A 함수의 런타임. 사용 가능한 값: java, python, go.
autoAck Boolean --auto-ack 프레임워크가 메시지를 자동으로 확인하는지 여부. 참고: 이 구성은 향후 릴리스에서 폐기될 예정이에요. 전달 시맨틱을 지정하면 프레임워크가 자동으로 메시지를 확인해요. 프레임워크가 메시지를 자동 확인하지 않게 하려면 processingGuaranteesMANUAL로 설정하세요.
maxMessageRetries Int --max-message-retries 메시지 처리를 포기하기 전까지 재시도하는 횟수.
deadLetterTopic String --dead-letter-topic 성공적으로 처리되지 않은 메시지를 저장하는 데 사용되는 토픽.
subName String --subs-name 필요한 경우 입력 토픽 소비자에 사용되는 Pulsar source 구독의 이름.
parallelism Int --parallelism 함수의 병렬성 계수, 즉 실행할 함수 인스턴스 수.
resources Resources N/A N/A
fqfn String --fqfn 함수의 완전 정규 함수 이름(FQFN).
windowConfig WindowConfig N/A N/A
timeoutMs Long --timeout-ms 메시지 타임아웃(밀리초). 참고: 이 값은 부정 확인된 메시지에 적용되는 재전달 지연도 설정해요. 함수는 이 둘을 독립적으로 구성할 방법이 없어요. sink에 해당하는 --negative-ack-redelivery-delay-ms에는 함수에 해당하는 것이 없어요. Java 함수는 둘 다 적용해요. Python 함수는 메시지 타임아웃만 적용하고 (#26411), Go 함수는 재전달 지연만 적용해요.
jar String --jar 함수(Java로 작성)의 JAR 파일 절대 경로. 워커가 패키지를 다운로드할 수 있는 URL 경로(HTTP, HTTPS, 현재 워커 호스트에 파일이 존재한다고 가정하는 file 프로토콜, 패키지 관리 서비스의 packages URL인 function)도 지원해요.
py String --py 함수(Python으로 작성)의 메인 Python/Python wheel 파일 절대 경로. 워커가 패키지를 다운로드할 수 있는 URL 경로(HTTP, HTTPS, file, function)도 지원해요.
go String --go 함수(Go로 작성)의 메인 Go 실행 바이너리 절대 경로. 워커가 패키지를 다운로드할 수 있는 URL 경로(HTTP, HTTPS, file, function)도 지원해요.
cleanupSubscription Boolean --cleanup-subscription 함수가 삭제될 때 함수가 만들거나 사용하는 구독을 삭제할지 여부. 기본값은 true.
customRuntimeOptions String --custom-runtime-options 런타임을 커스터마이즈하는 옵션을 인코딩한 문자열.
maxPendingAsyncRequests Int --max-message-retries 대량의 동시 요청을 피하기 위한 인스턴스당 보류 중인 비동기 요청의 최대 수.
exposePulsarAdminClientEnabled Boolean N/A Pulsar admin 클라이언트가 함수 컨텍스트에 노출되는지 여부. 기본적으로 비활성화돼 있어요.
subscriptionPosition String --subs-position 지정된 위치에서 메시지를 소비하는 데 사용되는 Pulsar source 구독의 위치. 기본값은 Latest.
skipToLatest Boolean --skip-to-latest 함수 인스턴스가 재시작될 때 소비자가 최신 메시지로 건너뛰어야 하는지 여부.
ConsumerConfig

다음 표는 inputSpecs 필드 아래의 중첩 필드와 관련 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
schemaType String N/A N/A
serdeClassName String N/A N/A
isRegexPattern Boolean N/A N/A
schemaProperties Map<String,String> N/A N/A
consumerProperties Map<String,String> N/A N/A
receiverQueueSize Int N/A N/A
cryptoConfig CryptoConfig N/A 코드 참조.
poolMessages Boolean N/A N/A
ProducerConfig

다음 표는 producerConfig 필드 아래의 중첩 필드와 관련 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
maxPendingMessages Int N/A 브로커로부터 확인을 받기 위해 대기 중인 메시지를 담는 큐의 최대 크기.
maxPendingMessagesAcrossPartitions Int N/A 모든 파티션에 걸친 maxPendingMessages 수.
useThreadLocalProducers Boolean N/A N/A
cryptoConfig CryptoConfig N/A 코드 참조.
batchBuilder String --batch-builder 배치 구성 방식의 유형. 사용 가능한 값: DEFAULTKEY_BASED. 기본값은 DEFAULT. 참고: batchingConfig.batchBuilder도 설정되어 있으면 그 값이 이 값보다 우선해요.
compressionType String N/A 프로듀서가 사용하는 메시지 데이터 압축 유형. 기본값은 LZ4. 사용 가능한 옵션: NONE(압축 없음), ZLIB, ZSTD, SNAPPY. 참고: Go 클라이언트는 SNAPPY를 구현하지 않아요. Go 함수가 SNAPPY로 구성되면 LZ4로 대체돼요.
batchingConfig BatchingConfig N/A 프로듀서에 적용되는 배칭 설정. 생략하면 최대 게시 지연 10ms로 배칭이 활성화돼요.
Resources

다음 표는 resources 필드 아래의 중첩 필드와 관련 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
cpu double --cpu 함수 인스턴스당 할당할 CPU(코어) (Kubernetes 런타임 전용).
ram Long --ram 함수 인스턴스당 할당할 RAM(바이트) (process/Kubernetes 런타임 전용).
disk Long --disk 함수 인스턴스당 할당할 디스크(바이트) (Kubernetes 런타임 전용).
WindowConfig

다음 표는 windowConfig 필드 아래의 중첩 필드와 관련 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
windowLengthCount Int --window-length-count 윈도우당 메시지 수.
windowLengthDurationMs Long --window-length-duration-ms 윈도우당 시간 기간(밀리초).
slidingIntervalCount Int --sliding-interval-count 윈도우가 슬라이드된 후의 메시지 수.
slidingIntervalDurationMs Long --sliding-interval-duration-ms 윈도우가 슬라이드된 후의 시간 기간.
lateDataTopic String N/A N/A
maxLagMs Long N/A N/A
watermarkEmitIntervalMs Long N/A N/A
timestampExtractorClassName String N/A N/A
actualWindowFunctionClassName String N/A N/A
CryptoConfig

다음 표는 cryptoConfig 필드 아래의 중첩 필드와 관련 인자를 설명해요.

필드 이름 유형 관련 명령 인자 설명
cryptoKeyReaderClassName String N/A 코드 참조.
cryptoKeyReaderConfig Map<String, Object> N/A N/A
encryptionKeys String[] N/A N/A
producerCryptoFailureAction ProducerCryptoFailureAction N/A N/A
consumerCryptoFailureAction ConsumerCryptoFailureAction N/A N/A
BatchingConfig

다음 표는 producerConfigbatchingConfig 필드 아래의 중첩 필드를 설명해요. 이 설정은 PIP-401에서 도입됐어요.

참고 batchingConfig는 현재 Java 함수에만 적용돼요. Python과 Go 런타임은 최대 게시 지연 10ms로 배칭을 하드코드로 활성화하고 이 설정을 무시해요 (#26390, #26391).

필드 이름 유형 관련 명령 인자 설명
enabled Boolean N/A 프로듀서가 메시지를 배칭하는지 여부. 기본값은 true.
batchingMaxPublishDelayMs Int N/A 프로듀서가 배치를 보내기 전에 기다리는 최대 시간(밀리초). 기본값은 10.
batchingMaxMessages Int N/A 배치의 최대 메시지 수. 설정하지 않으면 클라이언트 기본값이 적용돼요.
batchingMaxBytes Int N/A 배치의 최대 크기(바이트). 설정하지 않으면 클라이언트 기본값이 적용돼요.
batchBuilder String N/A 배치 구성 방식의 유형. 사용 가능한 값: DEFAULTKEY_BASED. producerConfig.batchBuilder보다 우선해요.
roundRobinRouterBatchingPartitionSwitchFrequency Int N/A 키 없는 메시지에 대해 라운드로빈 라우터가 파티션을 전환하는 빈도. batchingMaxPublishDelayMs의 배수로 표현돼요. 파티션 전환 주기는 frequency * batchingMaxPublishDelayMs예요.

producerConfig를 설정하지 않았거나 batchingConfig가 없는 producerConfig를 설정한 함수는 최대 게시 지연 10ms로 배칭이 활성화돼요. 이것은 이 설정들이 구성 가능해지기 전에 함수가 가졌던 동작이며, Java, Python, Go 런타임에서 동일해요.

예시

다음 예시는 YAML 또는 JSON을 사용해 함수를 구성하는 방법을 보여줘요.

tenant: "public"
namespace: "default"
name: "config-file-function"
inputs:
  - "persistent://public/default/config-file-function-input-1"
  - "persistent://public/default/config-file-function-input-2"
output: "persistent://public/default/config-file-function-output"
jar: "function.jar"
parallelism: 1
resources:
  cpu: 8
  ram: 8589934592
autoAck: true
userConfig:
  foo: "bar"
{
  "tenant": "public",
  "namespace": "default",
  "name": "config-file-function",
  "inputs": [
    "persistent://public/default/config-file-function-input-1",
    "persistent://public/default/config-file-function-input-2"
  ],
  "output": "persistent://public/default/config-file-function-output",
  "jar": "function.jar",
  "parallelism": 1,
  "resources": {
    "cpu": 8,
    "ram": 8589934592
  },
  "autoAck": true,
  "userConfig": {
    "foo": "bar"
  }
}

더 알아보기 (Learn more)