新闻详情

新闻详情

首页 / 资讯中心 / 详情

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

发布时间:2026/9/25 23:05:32来源:尧图网络
Apache Pulsar Functions 快速入门实战:从本地运行到集群部署
消息队列后端流处理【免费下载链接】pulsarApache 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 FunctionString, 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 模式的函数在运行函数的终端按CtrlC即可。刚才发生了什么发布到输入主题persistent://public/default/exclamation-input的Hello world消息被本机运行的 exclamation 函数接收函数处理消息后得到Hello world!并将结果发布到输出主题persistent://public/default/exclamation-output即使 exclamation 函数未处于运行状态发布到输入主题的消息也会被 Pulsar 持久化存储于 Apache BookKeeper 中直到有消费者消费并确认ack该消息——这正是 Pulsar 持久化保证的体现。以 cluster 模式运行 Pulsar Functionlocal 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 函数示例集更多可直接运行的参考实现。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

Google提示工程PDF拆解:从零样本到自洽性投票的实战指南 2026/9/25 23:42:06

Google提示工程PDF拆解:从零样本到自洽性投票的实战指南

简介:这份《google提示工程.pdf》面向具备一定编程基础、希望深入掌握大语言模型交互技巧的开发者、数据科学家与机器学习工程师,系统讲解如何编写高质量提示词以提升模型输出的准确性与相关性。内容覆盖零样本、少样本、系统提示、角色提示、上下文提示…

阅读更多 →
Google提示工程PDF实战:从零样本到结构化输出的提示词工程指南 2026/9/25 23:42:05

Google提示工程PDF实战:从零样本到结构化输出的提示词工程指南

简介:这份《google提示工程.pdf》面向具备一定编程基础、希望深入掌握大语言模型交互技巧的开发者、数据科学家与机器学习工程师,系统讲解如何编写高质量提示词以提升模型输出的准确性与相关性。内容覆盖零样本、少样本、系统提示、角色提示、上下文提示…

阅读更多 →
Nikon SDK C#开发实战:单拍、连拍与LiveView视频流 2026/9/25 23:41:59

Nikon SDK C#开发实战:单拍、连拍与LiveView视频流

简介:本资源是一套基于尼康官方SDK的C#与VB.NET相机控制开发套件,面向摄影自动化开发者、工业视觉工程师及高校计算机视觉方向学习者,解决尼康相机通过桌面软件实现视频录制、连拍、单拍等远程控制的核心需求。压缩包共63个文件,含…

阅读更多 →
央国企AI+数智化转型:从报告到落地的工程实践与避坑指南 2026/9/25 23:41:39

央国企AI+数智化转型:从报告到落地的工程实践与避坑指南

简介:这份《2025央国企AI数智化转型研究报告》面向央国企管理者、数字化转型负责人及产业研究者,系统梳理AI与大数据在央国企落地中的战略路径、技术应用与生态协同问题。报告从发展现状、核心挑战与痛点切入,覆盖战略路径、技术数据、组织人…

阅读更多 →
基于 Java 的轻量级 AI Agent 平台开源:TaoToken 统一 Key 接入 Spring AI 实战 2026/9/25 23:41:20

基于 Java 的轻量级 AI Agent 平台开源:TaoToken 统一 Key 接入 Spring AI 实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
Spirula Studio vs 传统3DGS工具链:为什么开发者都在迁移到这个一体化平台 2026/9/25 23:41:20

Spirula Studio vs 传统3DGS工具链:为什么开发者都在迁移到这个一体化平台

Spirula Studio vs 传统3DGS工具链:为什么开发者都在迁移到这个一体化平台 【免费下载链接】spirula-studio Cross-vendor 3D Gaussian Splatting trainer - video to splat to mesh, Vulkan or CUDA. 项目地址: https://gitcode.com/GitHub_Trending/sp/spirula-…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉