- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本指南以 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 exclamationlocal 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
相关推荐
Apache Pulsar 2.3.0 快速上手:本地运行与集群部署 Pulsar Functions 实战指南
Apache Pulsar 2.3.0 快速上手:本地运行与集群部署 Pulsar Functions 实战指南 本文基于 Apache Pulsar 2.3.
消息队列后端流处理Apache Pulsar Functions 快速上手指南:从本地运行到集群部署的完整实践
Apache Pulsar Functions 快速上手指南:从本地运行到集群部署的完整实践 本篇技术指南以 Apache Pulsar 官方入门文档为基础,带
消息队列后端流处理Apache Pulsar Functions 部署实战:从本地运行到集群模式的完整指南
Apache Pulsar Functions 部署实战:从本地运行到集群模式的完整指南 本文基于 Apache Pulsar 官方文档 functions d
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考