news 2026/9/25 23:05:26

Apache Pulsar Functions 快速入门实战:从本地运行到集群部署

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar Functions 快速入门实战:从本地运行到集群部署
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

本指南以 Apache Pulsar 的 Pulsar Functions 轻量级流处理模型为主题,带你从零搭建一个 standalone 单机集群,先后以 local run 与 cluster 两种模式运行 Pulsar Function,并完成消息消费、函数触发、并行度调整与删除的完整生命周期管理。读完本文,你将掌握pulsar-admin functions系列命令的实战用法,并理解 Java 与 Python 两种函数 API 的编写与部署方式。

本文内容以 Pulsar Functions 快速入门文档 为主体,并结合作者仓库(版本基线约 2.10.x)中的函数 API 源码与示例函数进行原理补充。

前置条件

在跟随本教程操作之前,请确保你的机器上已经安装了 Apache Maven。Maven 主要用于从源码构建示例函数 JAR(若直接使用二进制发行包中预编译好的examples/api-examples.jar,则无需额外构建,但仍建议备好 Maven 以便自行编译示例)。

运行 standalone Pulsar 集群

Pulsar Functions 运行在 Pulsar 集群之上,因此第一步是先在本地启动一个 Pulsar 集群。最简单的做法是使用standalone模式——从术语表的定义可以看出,standalone 模式将集群所需的全部组件(broker、BookKeeper、ZooKeeper 等)运行在同一台机器上,非常适合开发与实验用途。

首先下载对应版本的二进制发行包并解压、启动:

$ wget pulsar:binary_release_url $ tar xvfz apache-pulsar-@pulsar:version@-bin.tar.gz $ cd apache-pulsar-@pulsar:version@ $ bin/pulsar standalone \ --advertised-address 127.0.0.1

说明:pulsar:binary_release_url与@pulsar:version@是文档模板占位符,实际使用时请替换为对应发布版本的下载地址与版本号(例如apache-pulsar-2.x.x-bin.tar.gz)。

standalone 启动后,public租户(tenant)与default命名空间(namespace)会自动创建。本教程后续所有示例均使用该租户与命名空间,主题的完整名称形如persistent://public/default/<topic>。

以 local run 模式运行 Pulsar Function

一个最简单的 Java 函数

首先从一个简单的 Java 函数开始:它从输入主题读取字符串消息,在字符串末尾追加一个感叹号,再发布到输出主题。仓库中对应的示例实现位于 ExclamationFunction.java,其核心代码如下:

package org.apache.pulsar.functions.api.examples; import java.util.function.Function; public class ExclamationFunction implements Function<String, String> { @Override public String apply(String input) { return String.format("%s!", input); } }

需要说明的是,仓库中存在两种实现形态:

  • 上例使用 JDK 标准库的java.util.function.Function接口,对应仓库中的 JavaNativeExclamationFunction.java,即“Java 原生函数”写法;
  • 另一份 ExclamationFunction.java 则实现了 Pulsar 自定义的 Function 接口,方法签名为O process(I input, Context context),并额外提供initialize(Context)与close()两个生命周期钩子(默认空实现),用于在实例启动时初始化资源、在停止时释放资源。

两种写法均被--classname参数所支持,功能等价,可根据是否需要在处理逻辑中访问Context(例如读取用户配置、记录日志、访问函数状态等)来选择。

包含上述函数及其他多个示例函数的 JAR 已随二进制发行包提供,位于解压目录的examples文件夹中(构建产物对应 Maven 模块pulsar-functions-api-examples,见 pulsar-functions/java-examples/pom.xml)。

通过 localrun 在集群外部运行函数

使用pulsar-admin functions localrun命令,可以在你的笔记本上(即 Pulsar 集群外部)运行该函数,它依然与集群通信、订阅输入主题并发布结果:

$ bin/pulsar-admin functions localrun \ --jar examples/api-examples.jar \ --classname org.apache.pulsar.functions.api.examples.ExclamationFunction \ --inputs persistent://public/default/exclamation-input \ --output persistent://public/default/exclamation-output \ --name exclamation

local run 模式要点:函数运行在集群之外(例如你的开发机),仅通过网络与集群交互,适合开发调试;该模式的详细说明见 Pulsar Functions 部署文档。

从源码实现看,localrun子命令定义在 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java:它把命令行参数序列化为FunctionConfigJSON,随后启动$PULSAR_HOME/bin/function-localrunner脚本并在本机进程中执行该函数;真正的本地运行器则由 pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java 实现。

支持多个输入主题

--inputs参数除了指定单个主题,也支持用逗号分隔的多个主题列表,例如:

--inputs topic1,topic2

函数会同时订阅这些输入主题,并将处理结果发布到--output指定的输出主题。

验证函数运行:消费输出主题

另开一个终端,使用pulsar-client工具订阅并监听输出主题:

$ bin/pulsar-client consume persistent://public/default/exclamation-output \ --subscription-name my-subscription \ --num-messages 0

--num-messages 0表示消费者将无限期监听该主题,而不是读取固定数量的消息后退出。

向输入主题生产消息

再开一个终端,向输入主题发布一条消息:

$ bin/pulsar-client produce persistent://public/default/exclamation-input \ --num-produce 1 \ --messages "Hello world"

此时在消费终端应能看到处理后的输出:

----- got message ----- Hello world!

至此,第一次函数运行成功。如需关闭 local run 模式的函数,在运行函数的终端按Ctrl+C即可。

刚才发生了什么

  • 发布到输入主题persistent://public/default/exclamation-input的Hello world消息,被本机运行的 exclamation 函数接收;
  • 函数处理消息后得到Hello world!,并将结果发布到输出主题persistent://public/default/exclamation-output;
  • 即使 exclamation 函数未处于运行状态,发布到输入主题的消息也会被 Pulsar 持久化存储于 Apache BookKeeper 中,直到有消费者消费并确认(ack)该消息——这正是 Pulsar 持久化保证的体现。

以 cluster 模式运行 Pulsar Function

local run 模式适合开发与实验,但在真实生产部署中,你需要让函数运行在cluster 模式下:函数运行于 Pulsar 集群内部,并由同一套pulsar-admin functions管理接口统一管理(命令参考见pulsar-admin functions)。

创建函数:create

下面的命令将之前本地运行的 exclamation 函数部署到集群内部:

$ bin/pulsar-admin functions create \ --jar examples/api-examples.jar \ --classname org.apache.pulsar.functions.api.examples.ExclamationFunction \ --inputs persistent://public/default/exclamation-input \ --output persistent://public/default/exclamation-output \ --name exclamation

成功后输出Created successfully。

查看函数列表:list

列出指定租户、命名空间下运行的所有函数:

$ bin/pulsar-admin functions list \ --tenant public \ --namespace default

此时列表中应只包含exclamation一个函数。

查看运行状态:getstatus

使用getstatus命令查看函数的运行状态:

$ bin/pulsar-admin functions getstatus \ --tenant public \ --namespace default \ --name exclamation

返回的 JSON 输出如下:

{ "functionStatusList": [ { "running": true, "instanceId": "0" } ] }

可以看到:该实例当前处于运行状态(running: true),且集群中运行着一个实例,其 ID 为 0。

查看函数配置详情:get

若想获取函数更完整的信息(主题、租户、命名空间、类名、并行度等),使用get命令:

$ bin/pulsar-admin functions get \ --tenant public \ --namespace default \ --name exclamation

返回的 JSON 输出如下:

{ "tenant": "public", "namespace": "default", "name": "exclamation", "className": "org.apache.pulsar.functions.api.examples.ExclamationFunction", "output": "persistent://public/default/exclamation-output", "autoAck": true, "inputs": [ "persistent://public/default/exclamation-input" ], "parallelism": 1 }

注意其中的parallelism: 1:该字段表示函数的并行实例数,目前只有一个实例在运行。

调整并行度:update

使用update命令将函数的并行度调整为 3,即让集群中同时运行 3 个函数实例:

$ bin/pulsar-admin functions update \ --jar examples/api-examples.jar \ --classname org.apache.pulsar.functions.api.examples.ExclamationFunction \ --inputs persistent://public/default/exclamation-input \ --output persistent://public/default/exclamation-output \ --tenant public \ --namespace default \ --name exclamation \ --parallelism 3

成功后输出Updated successfully。再次执行get命令,可以看到parallelism已变为 3:

{ "tenant": "public", "namespace": "default", "name": "exclamation", "className": "org.apache.pulsar.functions.api.examples.ExclamationFunction", "output": "persistent://public/default/exclamation-output", "autoAck": true, "inputs": [ "persistent://public/default/exclamation-input" ], "parallelism": 3 }

从CmdFunctions的源码结构可以推断,create、update、delete、list、get、getstatus(与status同义)等子命令均注册于 pulsar-admin functions 命令入口,并通过PulsarAdmin.functions()的 Admin API 与集群交互。parallelism提升后,多个函数实例会并行消费输入主题分区中的消息,从而提升吞吐能力。

删除函数:delete

最后,用delete命令关闭并移除运行中的函数:

$ bin/pulsar-admin functions delete \ --tenant public \ --namespace default \ --name exclamation

看到Deleted successfully输出,即表示你已经完整走通了 cluster 模式下函数的创建、更新与关闭流程。

编写并运行一个全新的函数:Python 版字符串反转

前面的示例使用并管理了一个预编译的 Java 函数。下面通过 Python API 从零编写一个自己的函数:同样接收字符串,但将字符串反转后发布到指定主题。

安装 Python 客户端依赖

在编写 Python 函数之前,需要先安装相关依赖:

$ pip install pulsar-client

编写函数代码

创建一个新的 Python 文件:

$ touch reverse.py

在文件中写入以下内容:

def process(input): return input[::-1]

这里的process方法定义了函数的处理逻辑:利用 Python 的切片语法将每个传入字符串反转。这是 Pulsar Functions 的 Python 简洁写法,只关注“输入→输出”的映射关系;仓库中对应的示例可以参考 native_exclamation_function.py(简单函数式写法)与 exclamation_function.py(继承pulsar.Function基类、带Context的类式写法)。

部署函数:create

在 cluster 模式下部署该 Python 函数,注意使用--py参数指定 Python 文件,--classname指定函数名(此处即文件名去掉.py后的reverse):

$ bin/pulsar-admin functions create \ --py reverse.py \ --classname reverse \ --inputs persistent://public/default/backwards \ --output persistent://public/default/forwards \ --tenant public \ --namespace default \ --name reverse

看到Created successfully后,函数即可开始接收消息。

触发函数:trigger

由于函数运行在 cluster 模式,可以使用trigger命令主动向函数发送一条消息并获取处理结果,无需额外启动生产/消费客户端:

$ bin/pulsar-admin functions trigger \ --name reverse \ --tenant public \ --namespace default \ --trigger-value "sdrawrof won si tub sdrawkcab saw gnirts sihT"

预期输出为:

This string was backwards but is now forwards

至此,你已成功编写一个全新的 Pulsar Function,以 cluster 模式部署到 standalone 集群中,并通过trigger命令验证了其正确性。

延伸阅读

  • Pulsar Functions API 详解:了解 Java 与 Python 函数 API 的完整能力(Context、用户配置、窗口函数等);
  • Pulsar Functions 部署指南:深入对比 local run 与 cluster 模式的部署细节与配置项;
  • pulsar-admin 命令参考:查看functions子命令的完整参数列表;
  • pulsar-client 命令行工具参考:了解生产、消费命令的更多选项;
  • 函数 API 核心接口源码:process/initialize/close的接口定义;
  • Java 函数示例集 与 Python 函数示例集:更多可直接运行的参考实现。
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 22:57:21

OpenClaw卸载残留清理指南:服务、配置、缓存三步彻底清除

卸载这类带后台服务的 AI 代理工具&#xff0c;最恼人的不是卸载本身&#xff0c;而是卸载完总觉得哪儿不对劲——端口还在监听&#xff0c;开机又弹出日志报错&#xff0c;翻遍系统目录还有一堆.json、.db、.log残留。OpenClaw 尤其典型&#xff0c;它既有 CLI 主程序&#xf…

作者头像 李华
网站建设 2026/9/25 22:53:54

杭州大平层全案整体设计服务商实力与用户口碑深度解析

什么是大平层全案整体设计大平层这类改善型住宅&#xff0c;拥有开阔的空间面积和优越的地段资源&#xff0c;已经成为众多改善型家庭的置业&#xff0c;而全案整体设计是适配大平层空间的专属家居服务模式&#xff0c;和传统家居服务有着本质区别。传统家居消费中&#xff0c;…

作者头像 李华
网站建设 2026/9/25 22:52:53

大宅设计公司避坑挑选指南:专业实力与用户口碑深度解析

大宅设计的底层逻辑&#xff1a;为什么你家的豪宅始终用不对空间说起大宅设计&#xff0c;很多人第一反应就是花钱买好看&#xff0c;但真正住过的业主都知道&#xff0c;一套能称之为家的大宅&#xff0c;从来不是效果图里的悬浮楼梯和网红软装堆砌出来的。从入户到起居&#…

作者头像 李华
网站建设 2026/9/25 22:51:42

基于Qt框架的幸存者游戏源码解析与改造实战

简介&#xff1a;这份基于Qt框架的幸存者游戏源码包&#xff0c;是南京大学高级程序设计课程的大作业&#xff0c;围绕C面向对象编程思想设计实现。项目包含基本地图与障碍物生成、玩家角色的移动/攻击/掉血/拾取、敌方单位移动策略与攻击逻辑、局内与全局双重强化系统、存档读…

作者头像 李华
网站建设 2026/9/25 22:40:58

BERTopic实战:从嵌入到聚类,轻松搞定语义主题建模

简介&#xff1a;这套BERTopic模型教程代码包面向自然语言处理与主题建模入门者&#xff0c;聚焦如何用BERT文本嵌入替换传统词袋表示&#xff0c;再经UMAP降维、HDBSCAN聚类和类间词频逆文档频率生成主题&#xff0c;解决传统LDA忽略语义关联的问题&#xff0c;适合文本挖掘、…

作者头像 李华