함수 관리

함수 관리 (Manage Functions)

Pulsar 함수(Function)는 토픽에서 메시지를 받아 처리하고 결과를 내보내는 계산 단위예요. 이 페이지에서는 함수를 생성·업데이트·시작·중지·재시작·나열·삭제하고, 정보·상태·통계를 조회하고, 트리거와 상태 저장을 하는 방법을 pulsar-admin CLI, REST API, Java admin API로 정리했어요.

출처: 문서

본문

tip

이 페이지는 자주 사용하는 일부 작업만 보여줘요. 최신·전체 정보는 아래 참고 문서를 확인하세요.

Category Method If you want to manage functions...
Pulsar CLI pulsar-admin, which lists all commands, flags, descriptions, and more. See the functions command
Pulsar admin APIs REST API, which lists all parameters, responses, samples, and more. See the /admin/v3/functions endpoint
Pulsar admin APIs Java admin API, which lists all classes, methods, descriptions, and more. See the functions method of the PulsarAdmin object

함수에 대해 다음 작업을 수행할 수 있어요.

함수 생성 (Create a function)

Admin CLI, REST API, Java admin API를 사용해 클러스터 모드(즉, Pulsar 클러스터에 배포)로 Pulsar 함수를 만들 수 있어요. 모든 인터페이스는 같은 두 가지 입력을 받아요.

  • 구성(configuration): FunctionConfig의 필드(tenant, namespace, name, className, inputs, output, parallelism, userConfig, resources, ...). 모든 필드는 모든 인터페이스에서 같은 이름으로 제공돼요 — CLI에서는 커맨드라인 옵션 또는 YAML 파일의 키로, REST API에서는 JSON functionConfig 부분의 키로, Java에서는 FunctionConfig 객체의 setter로요. 생성과 업데이트 모두 마찬가지예요.
  • 패키지(package), 다음 세 가지 형태 중 하나로:
Package How to pass it Notes
A file CLI: --jar, --py 또는 --go; REST: 데이터 파일 부분; Java: 파일 이름 인자 함수 워커에 업로드돼요
A URL CLI: 구성 파일의 jar: <url>; REST: url 폼 필드; Java: createFunctionWithUrl 함수 워커가 가져오며, URL은 아래 설명대로 허용되어야 해요
A built-in function 구성의 jar: builtin://<function name>, 패키지 없음 워커가 functionsDirectory에서 해당 함수를 사용해요

패키지 URL은 함수 워커가 가져오는 것이며 conf/functions_worker.yml의 구성으로 허용되어야 해요. 허용되지 않은 URL은 400 Function Package url is not valid:로 실패해요.

  • file:///path/on/the/worker: 경로가 워커의 functionsDirectory 안에 있어야 해요. enableReferencingFunctionsDirectoryFiles: true(기본값)일 때요.
  • http://... 또는 https://...: URL이 additionalEnabledFunctionsUrlPatterns의 정규식 중 하나와 일치해야 해요(기본적으로 비어 있음). functions 디렉터리 밖의 file:// 경로도 같은 방식으로 허용할 수 있어요.
  • function://tenant/namespace/name@version: 패키지 관리에 업로드된 패키지. functionsWorkerEnablePackageManagement: true가 필요해요.

Sources와 sinks는 대신 connectorsDirectory, enableReferencingConnectorDirectoryFiles, additionalEnabledConnectorUrlPatterns를 사용해요.

Admin CLI

create 하위 명령을 사용해요. 구성을 커맨드라인 옵션으로 줄 수 있어요.

pulsar-admin functions create \
    --tenant public \
    --namespace default \
    --name exclamation \
    --classname org.apache.pulsar.functions.api.examples.ExclamationFunction \
    --inputs persistent://public/default/test-input-topic \
    --output persistent://public/default/test-output-topic \
    --parallelism 1 \
    --jar $PWD/examples/api-examples.jar

또는 --function-config-file로 전달하는 YAML 파일에 보관할 수 있어요. 유지보수가 더 쉬운데, 버전 관리에 넣을 수 있고 나중에 update에 같은 파일을 재사용할 수 있거든요. 커맨드라인 옵션이 파일의 값을 덮어써요.

cat > exclamation.yaml <<EOF
tenant: public
namespace: default
name: exclamation
className: org.apache.pulsar.functions.api.examples.ExclamationFunction
inputs:
  - persistent://public/default/test-input-topic
output: persistent://public/default/test-output-topic
parallelism: 1
EOF
pulsar-admin functions create \
    --function-config-file exclamation.yaml \
    --jar $PWD/examples/api-examples.jar

패키지 URL이나 빌트인 함수라면 파일의 jar 키에 넣고(jar: https://... 또는 jar: builtin://<function name>), --jar는 생략해요.

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}

요청은 multipart/form-data POST로, JSON 구성은 functionConfig 부분(콘텐츠 타입 application/json으로 보내야 함)에, 패키지는 데이터 파일 부분 또는 url 폼 필드(-F "url=https://...")에 담겨요. curl로는 다음과 같이 해요.

cat > /tmp/functionconfig.json <<EOF
{
  "tenant": "public",
  "namespace": "default",
  "name": "exclamation",
  "className": "org.apache.pulsar.functions.api.examples.ExclamationFunction",
  "runtime": "JAVA",
  "inputs": ["persistent://public/default/test-input-topic"],
  "output": "persistent://public/default/test-output-topic",
  "parallelism": 1
}
EOF
curl -X POST \
  -H "Authorization: Bearer *** token)" \
  -F "functionConfig=@/tmp/functionconfig.json;type=application/json" \
  -F "data=@$PWD/examples/api-examples.jar;type=application/octet-stream" \
  http://localhost:8080/admin/v3/functions/public/default/exclamation

빌트인 함수라면 패키지를 보내지 않고 구성에 jar를 설정해요.

cat > /tmp/functionconfig.json <<EOF
{
  "tenant": "public",
  "namespace": "default",
  "name": "myfunction",
  "jar": "builtin://builtin-function-name",
  "runtime": "JAVA",
  "inputs": ["persistent://public/default/input-topic"],
  "output": "persistent://public/default/output-topic",
  "parallelism": 1
}
EOF
curl -X POST \
  -H "Authorization: Bearer *** token)" \
  -F "functionConfig=@/tmp/functionconfig.json;type=application/json" \
  http://localhost:8080/admin/v3/functions/public/default/myfunction

함수 워커가 브로커와 함께 실행될 때는 요청을 브로커의 웹 서비스 포트(8080)로 보내고, 별도로 실행될 때는 워커 자체 포트(workerPort, 기본 6750)로 보내요. Sources와 sinks는 sourceConfig 또는 sinkConfig 부분으로 같은 형태를 사용해요.

Java

FunctionConfig functionConfig = new FunctionConfig();
functionConfig.setTenant(tenant);
functionConfig.setNamespace(namespace);
functionConfig.setName(functionName);
functionConfig.setRuntime(FunctionConfig.Runtime.JAVA);
functionConfig.setParallelism(1);
functionConfig.setClassName("org.apache.pulsar.functions.api.examples.ExclamationFunction");
functionConfig.setProcessingGuarantees(FunctionConfig.ProcessingGuarantees.ATLEAST_ONCE);
functionConfig.setTopicsPattern(sourceTopicPattern);
functionConfig.setSubName(subscriptionName);
functionConfig.setOutput(sinkTopic);
admin.functions().createFunction(functionConfig, fileName);

함수 업데이트 (Update a function)

이미 배포된 함수를 Admin CLI, REST API, Java admin API로 업데이트할 수 있어요. 업데이트는 create와 같은 구성과 패키지를 받고 같은 요청을 사용해요. 위 예제가 그대로 적용돼요. 다른 점은 함수 워커가 이것들을 어떻게 다루는가예요.

  • 구성이 배포된 구성과 병합(merge)돼요. 생략한 설정은 현재 값을 유지하고, tenant, namespace, name은 배포된 함수와 일치해야 해요. 따라서 함수가 생성될 때의 전체 구성을 보내거나, 식별 정보(tenant, namespace, name)와 변경할 설정만 보낼 수 있어요.
  • 일부 설정은 업데이트로 변경할 수 없어요. 입력 토픽, 구독 이름, 처리 보장(processing guarantees), 순서 보장(ordering guarantees), 런타임이 그것이에요. 그것들을 바꾸려면 함수를 삭제하고 다시 만들어야 해요.
  • 패키지는 선택사항이에요. 생략하면 배포된 코드를 유지하고, 제공하면 새 빌드를 롤아웃해요.
  • 구성도 패키지도 변경하지 않는 업데이트는 400 Update contains no change로 거부돼요.
  • 업데이트 전용 옵션인 --update-auth-data(REST·Java API의 updateOptions.updateAuthData)는 워커가 함수에 저장된 인증 데이터, 예를 들어 함수가 Pulsar에 연결할 때 쓰는 토큰을 호출자의 자격 증명으로 교체하게 해요.

Admin CLI

update 하위 명령을 사용해요. 업데이트는 배포된 구성에 병합되므로, create 예제와 같은 구성 파일이나 tenant, namespace, name과 변경할 설정만 담긴 파일, 또는 커맨드라인 옵션만 전달할 수 있어요. 어느 쪽이든 새 구현도 롤아웃하려면 --jar(또는 --py, --go)를 넘기고, 배포된 것을 유지하려면 생략해요.

예제:

# roll out a new build of the function (with whatever the file contains, changed or not)
pulsar-admin functions update \
    --function-config-file exclamation.yaml \
    --jar $PWD/examples/api-examples-2.jar
# change only the configuration, for example after setting parallelism: 2 in the file;
# the deployed implementation is kept
pulsar-admin functions update \
    --function-config-file exclamation.yaml
# the same change with command-line options only: the function's identity and the settings to change
pulsar-admin functions update \
    --tenant public \
    --namespace default \
    --name exclamation \
    --parallelism 2

REST API: PUT /admin/v3/functions/{tenant}/{namespace}/{functionName}

create와 같은 multipart/form-data 요청을 PUT으로 보내요. 업데이트는 배포된 구성에 병합되므로 functionConfig 부분은 전체 구성이거나 함수의 식별 정보와 변경할 설정만일 수 있어요. 새 구현도 롤아웃하려면 data 부분(또는 url 필드)을 포함하고, 배포된 것을 유지하려면 생략하고, 저장된 인증 데이터를 교체하려면 선택적 updateOptions 부분을 추가해요.

# roll out a new build of the function (with whatever the file contains, changed or not)
curl -X PUT \
  -H "Authorization: Bearer *** token)" \
  -F "functionConfig=@/tmp/functionconfig.json;type=application/json" \
  -F "data=@$PWD/examples/api-examples-2.jar;type=application/octet-stream" \
  http://localhost:8080/admin/v3/functions/public/default/exclamation
# change only the configuration, for example after setting "parallelism": 2 in the file;
# the deployed implementation is kept. Also replace the stored authentication data with the caller's
curl -X PUT \
  -H "Authorization: Bearer *** token)" \
  -F "functionConfig=@/tmp/functionconfig.json;type=application/json" \
  -F 'updateOptions={"updateAuthData":true};type=application/json' \
  http://localhost:8080/admin/v3/functions/public/default/exclamation
# the same change with a minimal configuration: the function's identity and the settings to change
curl -X PUT \
  -H "Authorization: Bearer *** token)" \
  -F 'functionConfig={"tenant":"public","namespace":"default","name":"exclamation","parallelism":2};type=application/json' \
  http://localhost:8080/admin/v3/functions/public/default/exclamation

Java

// Update merges into the deployed configuration, so the function's identity and the settings
// to change are enough; the complete configuration from the create example works as well.
FunctionConfig update = new FunctionConfig();
update.setTenant(tenant);
update.setNamespace(namespace);
update.setName(functionName);
// roll out a new build of the function, keeping the deployed configuration
admin.functions().updateFunction(update, "/path/to/api-examples-2.jar", new UpdateOptionsImpl());
// change only the configuration, here the parallelism; the deployed implementation is kept
update.setParallelism(2);
admin.functions().updateFunction(update, null, new UpdateOptionsImpl());
// also replace the stored authentication data with the caller's credentials
UpdateOptionsImpl updateOptions = new UpdateOptionsImpl();
updateOptions.setUpdateAuthData(true);
admin.functions().updateFunction(update, null, updateOptions);

UpdateOptionsImpl(org.apache.pulsar.common.functions, pulsar-common에 있음)은 updateFunction이 받는 UpdateOptions 인터페이스를 구현해요.

함수 시작 (Start a function)

함수의 인스턴스 하나를 시작하거나 함수의 모든 인스턴스를 시작할 수 있어요.

함수 인스턴스 하나 시작 (Start an instance of a function)

instance-id로 중지된 함수 인스턴스를 Admin CLI, REST API, Java admin API로 시작할 수 있어요.

Admin CLI

start 하위 명령을 사용해요.

pulsar-admin functions start \
    --tenant public \
    --namespace default \
    --name (the name of Pulsar Functions) \
    --instance-id 1

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/{instanceId}/start

Java

admin.functions().startFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));

함수의 모든 인스턴스 시작 (Start all instances of a function)

중지된 함수 인스턴스를 모두 Admin CLI, REST API, Java admin API로 시작할 수 있어요.

Admin CLI

start 하위 명령을 사용해요.

예제:

pulsar-admin functions start \
    --tenant public \
    --namespace default \
    --name (the name of Pulsar Functions) \

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/start

Java

admin.functions().startFunction(tenant, namespace, functionName);

함수 중지 (Stop a function)

함수의 인스턴스 하나를 중지하거나 함수의 모든 인스턴스를 중지할 수 있어요.

함수 인스턴스 하나 중지 (Stop an instance of a function)

instance-id로 함수 인스턴스를 Admin CLI, REST API, Java admin API로 중지할 수 있어요.

Admin CLI

stop 하위 명령을 사용해요.

예제:

pulsar-admin functions stop \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--instance-id 1

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/{instanceId}/stop

Java

admin.functions().stopFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));

함수의 모든 인스턴스 중지 (Stop all instances of a function)

모든 함수 인스턴스를 Admin CLI, REST API, Java admin API로 중지할 수 있어요.

Admin CLI

stop 하위 명령을 사용해요.

예제:

pulsar-admin functions stop \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/stop

Java

admin.functions().stopFunction(tenant, namespace, functionName);

함수 재시작 (Restart a function)

함수의 인스턴스 하나를 재시작하거나 함수의 모든 인스턴스를 재시작할 수 있어요.

함수 인스턴스 하나 재시작 (Restart an instance of a function)

instance-id로 함수 인스턴스를 Admin CLI, REST API, Java admin API로 재시작해요.

Admin CLI

restart 하위 명령을 사용해요.

예제:

pulsar-admin functions restart \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--instance-id 1

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/{instanceId}/restart

Java

admin.functions().restartFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));

함수의 모든 인스턴스 재시작 (Restart all instances of a function)

모든 함수 인스턴스를 Admin CLI, REST API, Java admin API로 재시작할 수 있어요.

Admin CLI

restart 하위 명령을 사용해요.

예제:

pulsar-admin functions restart \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/restart

Java

admin.functions().restartFunction(tenant, namespace, functionName);

모든 함수 나열 (List all functions)

특정 tenant와 namespace에서 실행 중인 모든 Pulsar 함수를 Admin CLI, REST API, Java admin API로 나열할 수 있어요.

Admin CLI

list 하위 명령을 사용해요.

예제:

pulsar-admin functions list \
	--tenant public \
	--namespace default

REST API: GET /admin/v3/functions/{tenant}/{namespace}

Java

admin.functions().getFunctions(tenant, namespace);

함수 삭제 (Delete a function)

Pulsar 클러스터에서 실행 중인 Pulsar 함수를 Admin CLI, REST API, Java admin API로 삭제할 수 있어요.

Admin CLI

delete 하위 명령을 사용해요.

예제:

pulsar-admin functions delete \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: DELETE /admin/v3/functions/{tenant}/{namespace}/{functionName}

Java

admin.functions().deleteFunction(tenant, namespace, functionName);

함수 정보 가져오기 (Get info about a function)

클러스터 모드에서 실행 중인 Pulsar 함수에 대한 정보를 Admin CLI, REST API, Java admin API로 얻을 수 있어요.

Admin CLI

get 하위 명령을 사용해요.

예제:

pulsar-admin functions get \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}

Java

admin.functions().getFunction(tenant, namespace, functionName);

함수 상태 가져오기 (Get status of a function)

함수의 인스턴스 하나의 상태를 가져오거나 함수의 모든 인스턴스 상태를 가져올 수 있어요.

함수 인스턴스 하나의 상태 (Get status of an instance of a function)

instance-id로 Pulsar 함수 인스턴스의 현재 상태를 Admin CLI, REST API, Java admin API로 얻을 수 있어요.

Admin CLI

status 하위 명령을 사용해요.

예제:

pulsar-admin functions status \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--instance-id 1

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}/{instanceId}/status

Java

admin.functions().getFunctionStatus(tenant, namespace, functionName, Integer.parseInt(instanceId));

함수의 모든 인스턴스 상태 (Get status of all instances of a function)

Pulsar 함수 인스턴스의 현재 상태를 Admin CLI, REST API, Java admin API로 얻을 수 있어요.

Admin CLI

status 하위 명령을 사용해요.

예제:

pulsar-admin functions status \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}/status

Java

admin.functions().getFunctionStatus(tenant, namespace, functionName);

함수 통계 가져오기 (Get stats of a function)

함수의 인스턴스 하나의 통계를 가져오거나 함수의 모든 인스턴스 통계를 가져올 수 있어요.

함수 인스턴스 하나의 통계 (Get stats of an instance of a function)

instance-id로 Pulsar 함수 인스턴스의 현재 통계를 Admin CLI, REST API, Java admin API로 얻을 수 있어요.

Admin CLI

stats 하위 명령을 사용해요.

예제:

pulsar-admin functions stats \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--instance-id 1

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}/{instanceId}/stats

Java

admin.functions().getFunctionStats(tenant, namespace, functionName, Integer.parseInt(instanceId));

함수의 모든 인스턴스 통계 (Get stats of all instances of a function)

Pulsar 함수의 현재 통계를 Admin CLI, REST API, Java admin API로 얻을 수 있어요.

Admin CLI

stats 하위 명령을 사용해요.

예제:

pulsar-admin functions stats \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions)

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}/stats

Java

admin.functions().getFunctionStats(tenant, namespace, functionName);

함수 트리거 (Trigger a function)

제공된 값으로 특정 Pulsar 함수를 트리거할 수 있어요.

Admin CLI

trigger 하위 명령을 사용해요.

예제:

pulsar-admin functions trigger \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--topic (the name of input topic) \
	--trigger-value \"hello pulsar\"
	# or --trigger-file (the path of trigger file)

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/trigger

Java

admin.functions().triggerFunction(tenant, namespace, functionName, topic, triggerValue, triggerFile);

함수와 연관된 상태 저장 (Put state associated with a function)

Pulsar 함수와 연관된 상태를 저장할 수 있어요.

Admin CLI

putstate 하위 명령을 사용해요.

예제:

pulsar-admin functions putstate \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--state "{\"key\":\"pulsar\", \"stringValue\":\"hello pulsar\"}"

REST API: POST /admin/v3/functions/{tenant}/{namespace}/{functionName}/state/{key}

Java

TypeReference<FunctionState> typeRef = new TypeReference<FunctionState>() {};
FunctionState stateRepr = ObjectMapperFactory.getThreadLocal().readValue(state, typeRef);
admin.functions().putFunctionState(tenant, namespace, functionName, stateRepr);

함수와 연관된 상태 가져오기 (Fetch state associated with a function)

Pulsar 함수와 연관된 현재 상태를 가져올 수 있어요.

Admin CLI

querystate 하위 명령을 사용해요.

예제:

pulsar-admin functions querystate \
	--tenant public \
	--namespace default \
	--name (the name of Pulsar Functions) \
	--key (the key of state)

REST API: GET /admin/v3/functions/{tenant}/{namespace}/{functionName}/state/{key}

Java

admin.functions().getFunctionState(tenant, namespace, functionName, key);

더 알아보기 (Learn more)

  • 함수 관련 전체 명령은 pulsar-admin의 functions 명령 참고서를 확인해요.
  • REST API 엔드포인트 상세는 /admin/v3/functions 문서를 참고해요.
  • Java admin API의 functions 메서드는 PulsarAdmin 객체 문서에서 볼 수 있어요.
  • 클러스터에 함수를 배포하는 전체 워크플로가 궁금하다면 Functions 시작하기 문서를 참고해요.