- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文基于 Apache Pulsar 官方文档 functions-deploy.md 编写,结合仓库源码与配置文件补充实现细节。Apache Pulsar Functions 提供轻量级的"Lambda 风格"计算能力,允许以函数方式消费消息、处理后写入输出 topic。本文将完整讲解部署前置条件、
pulsar-admin functions命令行接口、默认参数推断规则、本地运行模式与集群模式的区别、并行度与资源分配、包管理服务集成,以及如何通过trigger命令实时触发函数,帮助读者掌握从开发机到生产集群的完整部署链路。
前置要求:部署 Pulsar Functions 前需要准备什么
要部署和管理 Pulsar Functions,首先必须有一个正在运行的 Pulsar 集群。根据使用场景,可以选择以下任一方式搭建:
- 在本地机器上运行 standalone 模式 集群;
- 在 Kubernetes、Amazon Web Services、裸机(bare metal)、DC/OS 等环境上部署集群。
如果运行的不是 standalone 集群,则需要获取集群的service URL。获取方式取决于集群的部署方式。
此外,如果要在部署后触发(trigger)Python 用户自定义函数,必须在所有运行 functions worker 的机器上安装 pulsar python client。这是 Python 函数实例能够连接 broker、消费与生产消息的前提。
命令行接口:pulsar-admin functions 概览
Pulsar Functions 的部署与管理全部通过pulsar-admin functions接口完成。该接口包含多个子命令,常用的有:
create:以集群模式部署函数;trigger:触发已部署的函数(见下文"触发 Pulsar Functions");list:列出已部署的函数;- 其他命令还包括
update、delete、get、getstatus、getstats、restart、stop、start、localrun、upload、download等。
从源码结构看,这些子命令都在 CmdFunctions.java 中定义,该类通过 JCommander 框架解析命令行参数,@Parameters(commandDescription = "Interface for managing Pulsar Functions...")声明了命令整体说明。
默认参数:不指定时系统如何推断
管理 Pulsar Functions 时,需要指定大量函数信息,包括 tenant、namespace、输入/输出 topic 等。但其中部分参数在未指定时会有默认值。下表列出了完整的默认值规则:
| 参数 | 默认值 |
|---|---|
| 函数名(Function name) | 可对类名取任意值(除 org、library 等类似类名外)。例如指定--classname org.example.MyFunction时,函数名为MyFunction |
| Tenant | 从输入 topic 名称推导。如果输入 topic 位于marketingtenant 下(即 topic 名形如persistent://marketing/{namespace}/{topicName}),则 tenant 为marketing |
| Namespace | 从输入 topic 名称推导。如果输入 topic 位于marketingtenant 的asianamespace 下(topic 名形如persistent://marketing/asia/{topicName}),则 namespace 为asia |
| 输出 topic(Output topic) | {输入 topic}-{函数名}-output。例如输入 topic 名为incoming、函数名为exclamation,则输出 topic 名为incoming-exclamation-output |
| 订阅类型(Subscription type) | 对于at-least-once和at-most-once处理保证,默认应用SHARED模式;对于effectively-once保证,则应用FAILOVER模式 |
| 处理保证(Processing guarantees) | ATLEAST_ONCE |
| Pulsar service URL | pulsar://localhost:6650 |
默认参数示例:create 命令的实际行为
以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 - Tenant:
public - Namespace:
default - 订阅类型:
SHARED - 处理保证:
ATLEAST_ONCE - Pulsar service URL:
pulsar://localhost:6650
源码视角:默认参数是如何推断的
默认值的推断逻辑在仓库源码中有明确实现。在 CmdFunctions.java 中,NamespaceCommand.processArguments()在未指定--tenant和--namespace时分别将其置为PUBLIC_TENANT(即public)和DEFAULT_NAMESPACE(即default)。随后的FunctionCommand还支持使用--fqfn(Fully Qualified Function Name,形如tenant/namespace/name)一次性指定三者,且禁止--fqfn与--tenant/--namespace/--name混用,否则会抛出运行时异常。
函数名的推断则由 Utils.java 中的inferMissingFunctionName完成:它按.分割类名,取最后一段作为函数名——例如org.example.MyFunction推断为MyFunction。若未提供 tenant/namespace,inferMissingTenant与inferMissingNamespace同样回退到public与default。
此外,validateFunctionConfigs 还会做完整性校验:Python 与 Java 函数必须指定--classname(Go 函数不需要);必须且只能指定--jar、--py、--go三者之一;本地文件必须真实存在,或为受支持的包 URL。这些校验保证了配置在提交给集群前就是合法的。
本地运行模式(Local Run Mode)
在本地运行(local run)模式下,函数运行在执行命令的机器上——可以是开发者的笔记本电脑,也可以是 AWS EC2 实例等。下面是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 集群。如果希望本地运行但连接到非本地集群,可以使用--broker-service-url指定不同的 broker URL:
$ bin/pulsar-admin functions localrun \ --broker-service-url pulsar://my-cluster-host:6650 \ # Other function parameters从源码实现看,LocalRunner在 CmdFunctions.java 中定义了--broker-service-url(同时保留了旧的驼峰写法--brokerServiceUrl以兼容历史脚本),还支持--web-service-url、--client-auth-plugin、--use-tls、--tls-trust-cert-path等连接参数,以及--runtime(仅对 Java 函数生效,可选THREAD或PROCESS)、--metrics-port-start等运行时参数。
集群模式(Cluster Mode)
当函数以集群(cluster)模式运行时,函数代码会被上传到 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更新集群模式下的函数
可以使用update命令更新以集群模式运行的函数。下面的命令将上文创建的函数的输入、输出 topic 进行了更新:
$ bin/pulsar-admin functions update \ --py myfunc.py \ --classname myfunc.SomeFunction \ --inputs persistent://public/default/new-input-topic \ --output persistent://public/default/new-output-topic并行度(Parallelism)
Pulsar Functions 以进程或线程形式运行,这些运行单元被称为实例(instance)。默认情况下,一个函数只运行单个实例。通过一条localrun命令只能运行函数的一个实例;如需运行多个实例,需要多次执行localrun命令。
创建函数时,可以指定函数的并行度(即要运行的实例数量),使用create命令的--parallelism标志:
$ bin/pulsar-admin functions create \ --parallelism 3 \ # Other function info也可以使用update接口调整已创建函数的并行度:
$ bin/pulsar-admin functions update \ --parallelism 5 \ # Other function如果通过 YAML 文件指定函数配置,则使用parallelism参数。以下是一个配置文件示例:
# function-config.yaml parallelism: 3 inputs: - persistent://public/default/input-1 output: persistent://public/default/output-1 # other parameters对应的更新命令为:
$ bin/pulsar-admin functions update \ --function-config-file function-config.yaml从源码看,--parallelism参数与--function-config-file(同时兼容旧参数--functionConfigFile)都在FunctionDetailsCommand中声明;当提供配置文件时,会通过CmdUtils.loadConfig将 YAML 反序列化为FunctionConfig,命令行中显式指定的参数(如--parallelism)随后会覆盖配置文件中的同名项。
函数实例资源分配
以集群模式运行 Pulsar Functions 时,可以为每个函数 实例 指定分配的资源:
| 资源 | 指定方式 | 运行时 |
|---|---|---|
| CPU | 核数 | Kubernetes |
| 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 内存。进行资源规划时,务必把并行度(实例数量)纳入计算。
对应源码中,--cpu、--ram、--disk参数分别被解析为Double、Long、Long类型,并封装进Resources对象(--ram、--disk均为字节单位,因此 8 GB 对应8589934592、10 GB 对应10737418240)。
使用 Package Management 服务管理函数包
包管理(Package Management)服务实现了包的版本管理,简化 Functions、Sinks、Sources 的升级与回滚流程。当同一个函数、Sink 或 Source 需要在不同 namespace 中复用时,可以将它们上传到一个公共的包管理系统中统一管理。
要使用 Package management 服务,需要先在集群中启用该服务,在broker.conf中设置以下属性:
注意:Package management 服务默认不启用。
enablePackagesManagement=true packagesManagementStorageProvider=org.apache.pulsar.packages.management.storage.bookkeeper.BookKeeperPackagesStorageProvider packagesReplicas=1 packagesManagementLedgerRootPath=/ledgers在仓库自带的 conf/broker.conf 中可以看到这些配置的真实默认值:enablePackagesManagement=false(默认关闭)、packagesManagementStorageProvider默认即指向 BookKeeper 存储实现、packagesReplicas=1、packagesManagementLedgerRootPath=/ledgers。注释还说明:使用BookKeeperPackagesStorageProvider时,可通过bookkeeper_前缀为 BookKeeper 客户端追加配置。
启用后,可以通过 上传包 将函数包上传到服务中,并获得对应的 包 URL。拿到可用的包 URL 后,即可在pulsar-admin functions create中把--jar、--py或--go设置为该包 URL 来创建函数。
这一点在源码中同样有印证:--jar、--py、--go参数的描述(CmdFunctions.java)明确指出它们除了支持本地路径,还支持http/https/file协议 URL,以及来自包管理服务的function协议包 URL,由 worker 负责下载包。
触发 Pulsar Functions(Trigger)
如果一个 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)以 本地运行模式 创建该函数:
$ 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 # Listen indefinitely接着触发函数:
$ 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命令中,只需指定函数的基本信息(tenant、namespace 和 name)。触发函数时,不需要知道函数的输入 topic。
小结
本文完整梳理了 Apache Pulsar Functions 的部署链路:从集群前置准备开始,介绍了pulsar-admin functions命令行接口的常用子命令与默认参数推断规则(含源码级证明),随后分别讲解了本地运行模式与集群模式的差异、update更新、并行度与 YAML 配置、按实例计量的资源分配、包管理服务集成,以及基于trigger的实时调试方法。部署时建议遵循以下要点:显式指定 tenant/namespace/name 以避免依赖推断;集群模式下务必把并行度计入资源预算;需要多 namespace 复用函数包时提前在broker.conf中启用包管理服务;本地调试 Python 函数前确认所有 functions worker 机器已安装 pulsar python client。
</|DSML|parameter> </|DSML|invoke> </|DSML|tool_calls>
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar Functions 快速上手指南:从本地运行到集群部署的完整实践
Apache Pulsar Functions 快速上手指南:从本地运行到集群部署的完整实践 本篇技术指南以 Apache Pulsar 官方入门文档为基础,带
消息队列后端流处理Google IMA SDK Web(HTML5)客户端广告插入完整集成指南
Google IMA SDK Web(HTML5)客户端广告插入完整集成指南 本指南基于 ima sdk web guide.md https://link.g
消息队列后端流处理Apache Pulsar 裸机多集群部署完整指南:从 ZooKeeper、BookKeeper 到 Broker 的实战部署
Apache Pulsar 裸机多集群部署完整指南:从 ZooKeeper、BookKeeper 到 Broker 的实战部署 导读 本文是基于当前 Apach
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考