尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Pulsar Functions 部署与管理:Local Run、Cluster Mode、并行度与 Trigger 机制详解

发布时间:2026/9/25 2:13:04

资讯中心
01
ARTICLE

Apache Pulsar Functions 部署与管理:Local Run、Cluster Mode、并行度与 Trigger 机制详解

Apache Pulsar Functions 部署与管理:Local Run、Cluster Mode、并行度与 Trigger 机制详解
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文围绕 Apache Pulsar 的官方文档《Deploying and managing Pulsar Functions》展开系统讲解 Pulsar Functions 的两种部署模式Local run mode 与 Cluster mode、pulsar-admin functions命令行工具、FQFN 命名规范、各参数的默认值推导规则、并行度parallelism与实例资源配置以及通过trigger命令对函数进行端到端测试的完整流程并结合仓库中的 CLI 实现与函数配置工具类源码帮助读者理解这些默认值与订阅类型背后真实的代码逻辑读完即可独立地在自己的 Pulsar 集群中部署、更新、调参和验证 Pulsar Functions。两种部署模式在 2.3.0 版本中Pulsar Functions 提供了两种部署模式模式说明Local run mode本地运行模式函数运行在你的本地环境中例如你的笔记本电脑Cluster mode集群模式函数运行在 Pulsar 集群内部与 Pulsar broker 部署在同一批机器上需要注意的是Pulsar Functions 在设计之初就考虑了可扩展性未来可能支持更多部署选项如果希望贡献新的部署模式文档建议先与 Pulsar 开发者社区devpulsar.apache.org取得联系。前置要求要部署和管理 Pulsar Functions前提是要有一个运行中的 Pulsar 集群可选方案包括在本机运行一个 standalone 集群将 Pulsar 集群部署到 Kubernetes、AWS、裸机、DC/OS 等平台。如果运行的是非 standalone 集群还需要获取集群的 service URL具体获取方式取决于集群的部署方式。此外如果计划部署并触发 Python 编写的用户自定义函数应预先安装 Pulsar 的 Python 客户端库。命令行接口pulsar-admin functionsPulsar Functions 的部署与管理统一通过pulsar-admin functions接口完成其中包含create在集群模式下部署函数trigger向函数发送测试值并获取输出list列出已部署的函数以及localrun、update、delete、get、get-stats、restart、stop、start、get-state、put-state、upload、download等其余子命令。从源码结构看这些子命令全部实现在 CmdFunctions.java 中它继承CmdBase内部按子命令拆分为LocalRunner、CreateFunction、DeleteFunction、UpdateFunction、GetFunction、GetFunctionStats、RestartFunction、StopFunction、StartFunction、ListFunctions、StateGetter、StatePutter、TriggerFunction、UploadFunction、DownloadFunction等内部类见 CmdFunctions.java 第 66~82 行。这意味着pulsar-admin functions下的每个子命令都对应一个带独立参数解析逻辑的 CLI 命令实现。Fully Qualified Function NameFQFN每个 Pulsar Function 都有一个完全限定函数名FQFN由三个元素组成租户tenant、命名空间namespace和函数名name形式为tenant/namespace/name借助 FQFN你可以创建多个同名函数只要它们位于不同的命名空间中即可。在 CLI 层面--fqfn与--tenant/--namespace/--name是二选一的关系源码中 FunctionCommand.processArguments() 会校验两者不能同时出现并将 FQFN 按/切分为恰好 3 段后分别填充 tenant、namespace、functionName若只使用三个独立参数而未指定名称则直接抛出“必须指定函数名或 FQFN”的异常。默认参数Default arguments管理 Pulsar Functions 时需要指定租户、命名空间、输入输出 topic 等信息但以下参数在省略时会使用默认值参数默认值Function name函数名取类名去掉 org、library 等前缀的值。例如--classname org.example.MyFunction会使函数名为MyFunctionTenant租户从输入 topic 名推导。若输入 topic 位于marketing租户下即 topic 形如persistent://marketing/{namespace}/{topicName}则租户为marketingNamespace命名空间从输入 topic 名推导。若输入 topic 位于marketing租户下的asia命名空间即 topic 形如persistent://marketing/asia/{topicName}则命名空间为asiaOutput topic输出 topic{输入 topic}-{函数名}-output。例如输入 topic 为incoming、函数名为exclamation时输出 topic 为incoming-exclamation-outputSubscription type订阅类型对于 at-least-once 和 at-most-once 的处理保证默认应用SHARED对于 effectively-once 保证默认应用FAILOVERProcessing guarantees处理保证ATLEAST_ONCE见 functions-guaranteesPulsar service URLpulsar://localhost:6650关于订阅类型的默认推导仓库源码给出了精确依据FunctionConfigUtils.java 第 149~159 行 中当开启retainOrdering或处理保证为EFFECTIVELY_ONCE时选择FAILOVER当开启retainKeyOrdering时选择KEY_SHARED否则默认落到SHARED。这与文档中的默认值表格完全吻合。默认值使用示例以这条create命令为例$ bin/pulsar-admin functions create \ --jar my-pulsar-functions.jar \ --classname org.example.MyFunction \ --inputs my-function-input-topic1,my-function-input-topic2创建出的函数将获得如下默认值函数名MyFunction、租户public、命名空间default、订阅类型SHARED、处理保证ATLEAST_ONCE、Pulsar service URLpulsar://localhost:6650。这里“租户默认public、命名空间默认default”在源码中可以印证FunctionCommand.processArguments() 在tenant null时赋值为PUBLIC_TENANTnamespace null时赋值为DEFAULT_NAMESPACE。函数配置本身由 FunctionConfig.java 定义包含tenant、namespace、name、className、output、processingGuarantees、parallelism、jar、py、go等字段见 FunctionConfig.java 第 58~125 行CLI 参数最终都会被组装进该对象再提交给 Functions worker。Local run mode本地运行模式在local run模式下函数运行在发出命令的那台机器上可以是你的笔记本、一台云实例等。localrun命令示例$ bin/pulsar-admin functions localrun \ --py myfunc.py \ --classname myfunc.SomeFunction \ --inputs persistent://public/default/input-1 \ --output persistent://public/default/output-1默认情况下函数会通过本机 broker 的 service URLpulsar://localhost:6650连接同机的 Pulsar 集群。如果想用 local run 模式运行函数但连接非本地的 Pulsar 集群可以用--brokerServiceUrl标志指定其他 broker URL$ bin/pulsar-admin functions localrun \ --broker-service-url pulsar://my-cluster-host:6650 \ # 其他函数参数对应源码中LocalRunner子命令在 CmdFunctions.java 第 690 行起 定义其中--broker-service-url参数描述为“The URL for Pulsar broker”第 703 行它决定了本地运行时函数连接的 broker 地址。Cluster mode集群模式在cluster mode下函数代码会被上传到一个 Pulsar broker并与 broker 一起运行而不是在你的本地环境中运行。使用create命令即可在集群模式运行函数$ bin/pulsar-admin functions create \ --py myfunc.py \ --classname myfunc.SomeFunction \ --inputs persistent://public/default/input-1 \ --output persistent://public/default/output-1集群模式下函数的实际承载者是 Functions worker其运行参数如 worker 标识、函数包副本数、函数分配 topic 等集中在 functions_worker.yml 中配置例如workerId、numFunctionPackageReplicas、functionAssignmentTopicName等字段可结合该文件了解 worker 侧的完整配置项。更新集群模式函数使用update命令可以更新处于集群模式的 Pulsar Function。例如更新上文创建的函数$ bin/pulsar-admin functions update \ --py myfunc.py \ --classname myfunc.SomeFunction \ --inputs persistent://public/default/new-input-topic \ --output persistent://public/default/new-output-topicParallelism并行度Pulsar Functions 以称为instance实例的进程形式运行。默认情况下函数以单个实例运行且在 local run 模式下只能运行单实例。创建函数时可以指定函数的并行度即要运行的实例数方法是使用create命令的--parallelism标志$ bin/pulsar-admin functions create \ --parallelism 3 \ # 其他函数信息也可以通过update接口调整已创建函数的并行度$ bin/pulsar-admin functions update \ --parallelism 5 \ # 其他函数参数--parallelism、--cpu、--ram、--disk等参数均在 CmdFunctions.java 第 310~319 行 中定义分别对应“函数并行度”“需要分配的 CPU 核数”“需要分配的内存字节数”“需要分配的磁盘字节数”。如果使用 YAML 指定函数配置则使用parallelism参数。示例配置文件# function-config.yaml parallelism: 3 inputs: - persistent://public/default/input-1 output: persistent://public/default/output-1 # 其他参数对应的更新命令$ bin/pulsar-admin functions update \ --function-config-file function-config.yaml--function-config-file参数在 CmdFunctions.java 第 275 行 定义允许整个函数配置以配置文件形式传入。函数实例资源Function instance resources在集群模式运行 Pulsar Functions 时可以为每个函数实例指定分配的资源资源指定方式支持的运行时CPU核数Docker即将支持RAM字节数Process、Docker磁盘空间字节数Docker示例为函数分配 8 核 CPU、8 GB 内存和 10 GB 磁盘空间$ bin/pulsar-admin functions create \ --jar target/my-functions.jar \ --classname org.example.functions.MyFunction \ --cpu 8 \ --ram 8589934592 \ --disk 10737418240资源是按实例计算的应用于某个 Pulsar Function 的资源是按函数的每个实例应用的。例如对一个并行度为 5 的函数应用 8 GB 内存实际上是为该函数总共应用了 40 GB 内存。做资源规划时务必把并行度即实例数量计入你的计算。Triggering Pulsar Functions触发函数如果 Pulsar Function 正在以集群模式运行可以随时通过命令行trigger触发它。触发函数意味着向函数发送一条带有特定值的消息并通过命令行拿到函数的输出如果有。触发函数本质上与在函数某个输入 topic 上生产一条消息来调用它没有区别。pulsar-admin functions trigger命令只是一个便捷机制让你无需使用pulsar-client工具或语言特定的客户端库即可向函数发送消息。下面用一个简单的 Python 函数演示完整的触发流程。该函数根据输入返回一个简单字符串# myfunc.py def process(input): return This function has been triggered with a value of {0}.format(input)先用create命令创建该函数$ bin/pulsar-admin functions create \ --tenant public \ --namespace default \ --name myfunc \ --py myfunc.py \ --classname myfunc \ --inputs persistent://public/default/in \ --output persistent://public/default/out然后使用pulsar-client consume命令在输出 topic 上启动一个消费者监听来自myfunc函数的消息$ bin/pulsar-client consume persistent://public/default/out \ --subscription-name my-subscription \ --num-messages 0 # 持续监听接着触发该函数$ bin/pulsar-admin functions trigger \ --tenant public \ --namespace default \ --name myfunc \ --trigger-value hello world监听输出 topic 的消费者随后会打印出----- got message ----- This function has been triggered with a value of hello world无需知道 topic 信息在上文的trigger命令中你可能注意到只需要指定函数的基本信息租户、命名空间、名称。要触发函数你并不需要了解该函数的输入 topic 是什么——TriggerFunction子命令只需 FQFN 与--trigger-value见 CmdFunctions.java 第 1053~1058 行具体投递到哪个输入 topic 由服务端根据已存储的函数配置完成。要点回顾操作命令关键说明本地运行pulsar-admin functions localrun单机运行--broker-service-url可指向远程集群集群模式创建pulsar-admin functions create函数上传至集群并与 broker 同机运行更新函数pulsar-admin functions update支持修改输入/输出、并行度也可用--function-config-file触发函数pulsar-admin functions trigger只需 FQFN --trigger-value无需知道输入 topic命名规范FQFN tenant/namespace/name支持跨命名空间同名函数结合仓库源码可以进一步深入的路径包括CLI 参数解析与子命令结构 CmdFunctions.java、函数配置模型 FunctionConfig.java、订阅类型与处理保证的推导逻辑 FunctionConfigUtils.java以及 worker 侧运行配置 functions_worker.yml。配合仓库中的函数示例脚本如run-exclamation-function.sh、run-counter-function.sh即可完整复现本文的部署与触发流程。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Functions 部署与管理实战本地运行、集群模式、并行度与触发机制Apache Pulsar Functions 部署与管理实战本地运行、集群模式、并行度与触发机制 导读 Pulsar Functions 是 Apache消息队列后端流处理Apache Pulsar Functions Worker 部署与管理实战与 Broker 合跑与独立运行两种模式详解Apache Pulsar Functions Worker 部署与管理实战与 Broker 合跑与独立运行两种模式详解 本文以 Apache Pulsar消息队列后端流处理CircularProgressDrawable 项目推荐CircularProgressDrawable 项目推荐 1. 项目基础介绍和主要编程语言 CircularProgressDrawable 是一个开源的 A消息队列后端流处理上一篇【亲测免费】 探索未来三维人像建模MMHuman3D下一篇探索高效数据库管理Sqlitedict - 简单、快速且本地化的SQLite接口创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。