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 |
프레임워크가 메시지를 자동으로 확인하는지 여부. 참고: 이 구성은 향후 릴리스에서 폐기될 예정이에요. 전달 시맨틱을 지정하면 프레임워크가 자동으로 메시지를 확인해요. 프레임워크가 메시지를 자동 확인하지 않게 하려면 processingGuarantees를 MANUAL로 설정하세요. |
| 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 |
배치 구성 방식의 유형. 사용 가능한 값: DEFAULT와 KEY_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
다음 표는 producerConfig의 batchingConfig 필드 아래의 중첩 필드를 설명해요.
이 설정은 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 | 배치 구성 방식의 유형. 사용 가능한 값: DEFAULT와 KEY_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)
- Functions 개념 — 함수의 기본 개념을 익혀요.
- pulsar-admin CLI 참조 — functions 명령을 자세히 다뤄요.
- Functions 배포 — 함수를 실제로 배포해요.
- 메시징 개념 — 메시지와 토픽의 기본을 익혀요.