Apache Pulsar Functions 部署实战:从本地运行到集群模式的完整指南
2026/9/23 17:18:22 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

本文基于 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:列出已部署的函数;
  • 其他命令还包括updatedeletegetgetstatusgetstatsrestartstopstartlocalrunuploaddownload等。

从源码结构看,这些子命令都在 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-onceat-most-once处理保证,默认应用SHARED模式;对于effectively-once保证,则应用FAILOVER模式
处理保证(Processing guarantees)ATLEAST_ONCE
Pulsar service URLpulsar://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,inferMissingTenantinferMissingNamespace同样回退到publicdefault

此外,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 函数生效,可选THREADPROCESS)、--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参数分别被解析为DoubleLongLong类型,并封装进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=1packagesManagementLedgerRootPath=/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

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询