Apache Pulsar Functions 全生命周期管理实战:pulsar-admin CLI、REST API 与 Java Admin API 完全指南
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文以 Apache Pulsar 的 Functions 管理为主题系统讲解如何在集群模式下通过pulsar-admin命令行、/admin/v3/functionsREST 接口和PulsarAdminJava Admin API 三种途径完成 Pulsar Functions 的创建、更新、启停、重启、查询、删除、触发与状态管理。读完本文你将掌握 Pulsar Functions 运维的完整命令集、核心配置参数语义以及底层实现原理可直接用于生产环境的函数生命周期管理。Pulsar Functions 是什么Pulsar Functions是 Apache Pulsar 内置的轻量级计算进程lightweight compute process其核心能力是从一个或多个 Pulsar topic 中消费消息对每条消息应用用户提供的处理逻辑将计算结果发布到另一个 topic。从编程模型上看它类似于 Lambda 风格的函数式处理无需部署独立的流处理集群如 Storm、Flink即可在 Pulsar 内部完成实时计算。仓库中site2/docs/functions-overview.md对该模型有完整介绍本文聚焦的是它的管理面如何对已部署的 Function 做全生命周期操作。一个最经典的示例是仓库自带的ExclamationFunction源码见 ExclamationFunction.java它对输入字符串追加一个感叹号后输出public class ExclamationFunction implements FunctionString, String { Override public String process(String input, Context context) { return String.format(%s!, input); } }这个类被打包在示例 jar 中下文所有管理操作都以它为例。三种管理方式总览Functions 可以通过以下三种方式管理管理方式说明Admin CLIpulsar-admin工具中的functions命令族命令行入口实现见 CmdFunctions.javaREST API/admin/v3/functions端点broker 侧资源类见 Functions.javaworker 侧实现见 FunctionsImplV2.javaJava Admin APIPulsarAdmin对象上的functions()方法接口定义见 Functions.javaJava 客户端入门可参考 client-libraries-java.md三种方式操作的是同一套管理语义只是入口不同。每种操作在下面的章节中都会给出三种写法。前置知识租户、命名空间与 FQFN绝大多数 Functions 管理命令都需要定位函数所在的tenant租户与namespace命名空间。从 CmdFunctions.java 的源码可以看到CLI 对这两个参数有默认值逻辑未显式指定时--tenant默认为public--namespace默认为default对应TopicName.PUBLIC_TENANT与TopicName.DEFAULT_NAMESPACE。除了用--tenant/--namespace/--name三件套定位函数CLI 还支持使用FQFNFully Qualified Function Name。FQFN 的格式必须是tenant/namespace/name三段式。源码中parseFullyQualifiedFunctionNameCmdFunctions.java会将其拆分为三个部分若与三件套同时使用或格式不是三段CLI 会直接报错。# 三件套写法 $ pulsar-admin functions get --tenant public --namespace default --name my-func # 等价 FQFN 写法 $ pulsar-admin functions get --fqfn public/default/my-func创建 FunctionCreate a function在集群模式cluster mode下创建 Pulsar Function即将其部署到 Pulsar 集群上执行。Admin CLI使用create子命令$ pulsar-admin functions create \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --inputs test-input-topic \ --output persistent://public/default/test-output-topic \ --classname org.apache.pulsar.functions.api.examples.ExclamationFunction \ --jar /examples/api-examples.jarREST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName请求体为 multipart 表单包含data函数代码包、url可选函数包下载地址与functionDetailsJSON 格式的 FunctionConfig。服务端入口在 broker 的 Functions.java。Java Admin APIFunctionConfig 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.setAutoAck(true); functionConfig.setOutput(sinkTopic); admin.functions().createFunction(functionConfig, fileName);create 子命令的完整参数表以下参数表来自 reference-pulsar-admin.md 的create命令说明同时与 CmdFunctions.java 中的字段一一对应Flag说明默认值--cpu每个函数实例分配的 CPU核数仅 Docker runtime 生效--ram每个函数实例分配的内存字节process/docker runtime 生效--disk每个函数实例分配的磁盘字节仅 Docker runtime 生效--auto-ack框架是否自动确认ack消息--subs-name输入 topic 消费者使用的订阅名称--classname函数的类名--custom-serde-inputs输入 topic 到 SerDe 类名的映射JSON 字符串--custom-schema-inputs输入 topic 到 Schema 属性的映射JSON 字符串--function-config-file指定 YAML 配置文件加载函数配置--inputs输入 topic 或 topic 列表多个以逗号分隔--log-topic函数日志写入的 topic--jarJava 函数 jar 路径也支持 http/https/file/function 包管理 URL--name函数名--namespace函数所属命名空间default--output输出 topic不指定则无输出--output-serde-classname输出消息使用的 SerDe 类--parallelism函数并行度即运行的实例数1--processing-guarantees处理保证投递语义可选[ATLEAST_ONCE, ATMOST_ONCE, EFFECTIVELY_ONCE]ATLEAST_ONCE--pyPython 函数主文件/Python Wheel 路径同样支持 URL--goGo 函数可执行二进制路径同样支持 URL--schema-type输出消息使用的内置 schema 类型或自定义 Schema 类名--sliding-interval-count窗口滑动的消息数--sliding-interval-duration-ms窗口滑动的时间间隔毫秒--tenant函数所属租户public--topics-pattern按命名空间下 topic 名称模式消费与--inputs互斥--user-config用户自定义配置键值对--window-length-count每个窗口的消息数--window-length-duration-ms窗口时间长度毫秒--dead-letter-topic处理失败消息发送到的死信 topic--fqfn函数的完全限定名tenant/namespace/name--max-message-retries消息处理失败的最大重试次数--retain-ordering函数是否按顺序消费和处理消息--retain-key-ordering是否按消息 key 保证同一实例处理同一 key--timeout-ms消息处理超时毫秒--producer-config自定义生产者配置JSON 字符串底层实现细节从源码看create命令的核心逻辑在CreateFunction.runCmd()CmdFunctions.java它会判断 jar/py/go 包路径是否为受支持的 URLhttp/https/file/function协议若是则调用createFunctionWithUrl否则走本地上传的createFunction。这解释了为什么--jar既能传本地路径又能传远程 URL。Java 侧FunctionConfig中与创建强相关的枚举定义在 FunctionConfig.javaProcessingGuaranteesATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCERuntimeJAVA、PYTHON、GO即 Pulsar Functions 支持三种运行时语言。更新 FunctionUpdate a function更新一个已部署到集群的 Pulsar Function。Admin CLI$ pulsar-admin functions update \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --output persistent://public/default/update-output-topic \ # other optionsupdate复用了与create相同的参数体系还额外支持--update-auth-data标志用于指示是否同步更新函数的鉴权数据。REST APIPUT /admin/v3/functions/:tenant/:namespace/:functionNameJava Admin APIFunctionConfig 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); UpdateOptions updateOptions new UpdateOptions(); updateOptions.setUpdateAuthData(updateAuthData); admin.functions().updateFunction(functionConfig, userCodeFile, updateOptions);从 Functions.java 接口可以看到updateFunction与createFunction一样同时提供了带UpdateOptions的重载版本以及updateFunctionWithUrl支持通过 URL 更新代码包和全部异步变体方便不同场景选用。启动 Function 实例StartPulsar Function 按parallelism参数可运行多个实例instance实例 ID 从 0 开始编号。以下命令用于启动已停止的实例。启动单个实例Admin CLI$ pulsar-admin functions start \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --instance-id 1REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/:instanceId/startJava Admin APIadmin.functions().startFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));启动全部实例不传--instance-id即启动该函数所有已停止的实例。Admin CLI$ pulsar-admin functions start \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/startJava Admin APIadmin.functions().startFunction(tenant, namespace, functionName);从 CmdFunctions.java 的StartFunction实现可以看到带--instance-id时走单实例重载否则走全部实例重载instance-id必须为数字否则报错提示。停止 Function 实例Stop停止单个实例Admin CLI$ pulsar-admin functions stop \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --instance-id 1REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/:instanceId/stopJava Admin APIadmin.functions().stopFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));停止全部实例Admin CLI$ pulsar-admin functions stop \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/stopJava Admin APIadmin.functions().stopFunction(tenant, namespace, functionName);重启 Function 实例Restart重启单个实例Admin CLI$ pulsar-admin functions restart \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --instance-id 1REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/:instanceId/restartJava Admin APIadmin.functions().restartFunction(tenant, namespace, functionName, Integer.parseInt(instanceId));重启全部实例Admin CLI$ pulsar-admin functions restart \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/restartJava Admin APIadmin.functions().restartFunction(tenant, namespace, functionName);列出所有 FunctionList列出指定租户和命名空间下运行的所有 Pulsar Functions。Admin CLI$ pulsar-admin functions list \ --tenant public \ --namespace defaultREST APIGET /admin/v3/functions/:tenant/:namespaceJava Admin APIadmin.functions().getFunctions(tenant, namespace);getFunctions返回函数名列表如[f1, f2, f3]其签名与异步版本见 Functions.java。删除 FunctionDelete删除运行在 Pulsar 集群上的 Pulsar Function。Admin CLI$ pulsar-admin functions delete \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIDELETE /admin/v3/functions/:tenant/:namespace/:functionNameJava Admin APIadmin.functions().deleteFunction(tenant, namespace, functionName);注意FunctionConfig中有一个cleanupSubscription字段见 FunctionConfig.java用于控制删除函数时是否一并清理其创建/使用的订阅。获取 Function 信息Get获取当前以集群模式运行的 Pulsar Function 的配置信息。Admin CLI$ pulsar-admin functions get \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIGET /admin/v3/functions/:tenant/:namespace/:functionNameJava Admin APIadmin.functions().getFunction(tenant, namespace, functionName);CLI 的GetFunction实现CmdFunctions.java会以 pretty-print 的 JSON 形式输出完整的FunctionConfig方便核对函数的全部当前配置。查看实例状态Status查看单个实例状态Admin CLI$ pulsar-admin functions status \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --instance-id 1REST APIGET /admin/v3/functions/:tenant/:namespace/:functionName/:instanceId/statusJava Admin APIadmin.functions().getFunctionStatus(tenant, namespace, functionName, Integer.parseInt(instanceId));查看全部实例状态Admin CLI$ pulsar-admin functions status \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIGET /admin/v3/functions/:tenant/:namespace/:functionName/statusJava Admin APIadmin.functions().getFunctionStatus(tenant, namespace, functionName);状态信息包含每个实例的运行健康度、运行时长等指标。status命令在 CLI 中还注册了getstatus别名见 CmdFunctions.java两个名字等价。查看实例统计Stats查看单个实例统计Admin CLI$ pulsar-admin functions stats \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --instance-id 1REST APIGET /admin/v3/functions/:tenant/:namespace/:functionName/:instanceId/statsJava Admin APIadmin.functions().getFunctionStats(tenant, namespace, functionName, Integer.parseInt(instanceId));查看全部实例统计Admin CLI$ pulsar-admin functions stats \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions)REST APIGET /admin/v3/functions/:tenant/:namespace/:functionName/statsJava Admin APIadmin.functions().getFunctionStats(tenant, namespace, functionName);stats与status的区别在于status关注实例的运行状态存活、健康stats关注吞吐与性能指标如接收/处理/输出的消息数。接口层面分别对应FunctionStatus与FunctionStats见 Functions.java。getFunctionStats命令同时被pulsar-admin functions-worker function-stats等运维工具复用。触发 FunctionTriggertrigger用于向函数的输入 topic 注入一条指定的数据从而触发函数执行一次适合功能验证与调试。Admin CLI$ 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)--trigger-value直接给出要注入的字符串--trigger-file则指定包含触发数据的文件路径二者至少提供其一。REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/triggerJava Admin APIadmin.functions().triggerFunction(tenant, namespace, functionName, topic, triggerValue, triggerFile);从 Functions.java 的注释与 CmdFunctions.java 的实现可以看出trigger的本质是向函数的输入 topic 写入一条消息值来自--trigger-value或--trigger-file由函数消费后执行处理逻辑触发成功后会返回函数处理的结果因此非常适合在部署后快速验证函数逻辑是否符合预期。--topic用于在函数有多个输入 topic 时指定注入目标。状态存储写入与读取Putstate / QuerystatePulsar Functions 支持为函数关联状态存储state storage函数可通过Context读写键值状态而管理员也可以从管理面直接写入或查询函数状态。写入状态PutstateAdmin CLI$ pulsar-admin functions putstate \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --state {\key\:\pulsar\, \stringValue\:\hello pulsar\}REST APIPOST /admin/v3/functions/:tenant/:namespace/:functionName/state/:keyJava Admin APITypeReferenceFunctionState typeRef new TypeReferenceFunctionState() {}; FunctionState stateRepr ObjectMapperFactory.getThreadLocal().readValue(state, typeRef); admin.functions().putFunctionState(tenant, namespace, functionName, stateRepr);查询状态QuerystateAdmin CLI$ pulsar-admin functions querystate \ --tenant public \ --namespace default \ --name (the name of Pulsar Functions) \ --key (the key of state)REST APIGET /admin/v3/functions/:tenant/:namespace/:functionName/state/:keyJava Admin APIadmin.functions().getFunctionState(tenant, namespace, functionName, key);FunctionState 数据结构状态对象的结构定义在 FunctionState.java包含以下字段字段说明key状态的键stringValue字符串形式的值byteValue字节数组形式的值numberValue数值形式的值version状态的版本号用于乐观并发控制在 CmdFunctions.java 中querystate还支持-w/--watch参数启用后命令会以 1 秒为周期持续轮询该 key 的状态值变化默认关闭非常适合在函数运行过程中观察状态更新。同时注意querystate要求必须指定--key否则报错State key needs to be specified。其他常用命令除了上述与原文对应的核心操作pulsar-admin functions命令族还包含以下实用子命令均注册于 CmdFunctions.java子命令说明localrun在本地运行 Pulsar Function而非部署到集群支持--broker-service-url、--web-service-url、TLS 与鉴权相关参数upload将函数代码包上传到 Pulsar配合包管理服务使用download从 Pulsar 下载函数代码包可通过包路径或 tenant/namespace/name 定位其中localrun的实现会拼接PULSAR_HOME/bin/function-localrunner命令并以进程方式启动本地运行器CmdFunctions.java适合在开发阶段快速验证函数逻辑后再提交到集群。生产实践要点明确租户/命名空间虽然 CLI 对--tenantpublic与--namespacedefault有默认值但在多租户生产环境建议始终显式指定避免误操作到默认命名空间下的函数。善用--function-config-file当参数较多时可将完整FunctionConfig写成 YAML 文件再通过该参数加载比超长命令行更易维护CLI 内部通过CmdUtils.loadConfig将 YAML 反序列化为FunctionConfigCmdFunctions.java。正确处理投递语义--processing-guarantees默认ATLEAST_ONCE如需要去重等更强语义可选择EFFECTIVELY_ONCE结合--auto-ack、--max-message-retries与--dead-letter-topic可构建完整的失败重试与死信链路。用trigger做冒烟测试部署后先用trigger注入一条测试数据并观察返回结果确认函数逻辑与输入输出 topic 配置无误后再接入真实流量。用status/stats做健康巡检日常巡检优先看status判断实例是否健康性能分析再深入statsquerystate -w可在排障时实时观察有状态函数的内部状态变化。以上所有操作命令与参数均以当前仓库代码为准CLI 实现见 CmdFunctions.javaJava Admin API 接口见 Functions.java配置模型见 FunctionConfig.javaworker 侧集群配置参考 functions_worker.yml。命令的完整 flag 清单可进一步查阅 reference-pulsar-admin.md 的functions章节。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐3步构建个人数字图书馆novel-downloader的跨平台内容聚合解决方案3步构建个人数字图书馆novel downloader的跨平台内容聚合解决方案 在数字阅读时代我们每天都在产生和消耗海量内容但真正的知识资产却常常流离失所消息队列后端流处理Apache Pulsar Admin 接口完全指南pulsar-admin CLI、REST API 与 Java Admin APIApache Pulsar Admin 接口完全指南pulsar admin CLI、REST API 与 Java Admin API 导读 Apache消息队列后端流处理Apache Pulsar Namespace 管理实战pulsar-admin、REST API 与 Java Admin API 全指南Apache Pulsar Namespace 管理实战pulsar admin、REST API 与 Java Admin API 全指南 导读 命名空间消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考