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),仅供参考

相关新闻

8B模型LoRA微调营销文案:数据准备、训练参数与Ollama部署

8B模型LoRA微调营销文案:数据准备、训练参数与Ollama部署

简介:面向具备机器学习基础的技术人员与市场营销从业者,一套围绕AI模型高效训练的实战指南,核心思路是先借助大型模型生成多样化营销训练数据,再通过Unsloth微调8B小模型,使其在广告文案、社交话题等营销内容生成上接近…

2026/9/25 23:05:38 阅读更多 →
YOLOv8接入RTSP流实时目标检测:拉流、避坑与低延迟实践

YOLOv8接入RTSP流实时目标检测:拉流、避坑与低延迟实践

简介:面向需要构建实时视频分析系统的开发者,这份YOLOv8基于RTSP流的目标检测资源包,提供了从视频流接入、图像预处理、模型推理到结果可视化的完整可运行方案,并覆盖环境配置与部署运行的关键细节。借助YOLOAPI工具与YAML配置文件…

2026/9/25 23:04:37 阅读更多 →
湖北煤矿道岔,双开道岔,盾构道岔,单开道岔优质厂家实力参考:林州市创扬矿山设备制造有限公司靠谱定制厂家推荐

湖北煤矿道岔,双开道岔,盾构道岔,单开道岔优质厂家实力参考:林州市创扬矿山设备制造有限公司靠谱定制厂家推荐

湖北地区煤矿企业轨道改造、新建矿井轨道铺设,不少采购负责人都在打听靠谱的煤矿道岔专业制造商,想找能做煤矿道岔个性化定制厂家,不少人问起推荐一下煤矿道岔制造商哪家靠谱,今天我们就结合行业实际情况,给大家聊一聊…

2026/9/25 23:04:37 阅读更多 →

最新新闻

Python f-string性能原理与工程实践指南

Python f-string性能原理与工程实践指南

1. 为什么我彻底停用了.format()和%,只用 f-string?三年前我还在带一个刚转行的实习生,他写了一段爬虫日志记录代码:log_msg "Request to {url} failed with status {code}, retrying {count} times".format(urlendpoi…

2026/9/27 0:51:05 阅读更多 →
Agentic调度系统:面向智能体的原生Kubernetes编排方案

Agentic调度系统:面向智能体的原生Kubernetes编排方案

1. 项目概述:从“ax”这个极简标题切入,我们到底在谈什么?“ax”——两个字母,没有空格,没有标点,没有上下文。乍一看像缩写、像代号、像密码,甚至像打字错误。但结合当前技术社区高频出现的热搜…

2026/9/27 0:51:05 阅读更多 →
如何开发FluentTweaker扩展:元数据、Host模式与参数传递完整参考

如何开发FluentTweaker扩展:元数据、Host模式与参数传递完整参考

如何开发FluentTweaker扩展:元数据、Host模式与参数传递完整参考 【免费下载链接】FluentTweaker Windows Slop Remover 项目地址: https://gitcode.com/gh_mirrors/wi/FluentTweaker FluentTweaker(Windows Slop Remover)是一款开源的…

2026/9/27 0:51:05 阅读更多 →
Agent Cookie Sync:Grok Bot与Muse的Chrome会话同步实践

Agent Cookie Sync:Grok Bot与Muse的Chrome会话同步实践

1. 从"登录态丢失"说起:Agent Cookie Sync 到底在解决什么做过浏览器自动化的人大概率都遇到过这个场景:脚本跑得好好的,突然某一天所有请求全部返回未登录,页面跳回登录页,之前辛苦维持的会话状态一夜清零。…

2026/9/27 0:51:05 阅读更多 →
企业高端网站制作避坑指南:5个技术选型坑,让网站真正被搜到

企业高端网站制作避坑指南:5个技术选型坑,让网站真正被搜到

企业高端网站制作避坑指南:5个技术选型坑,让网站真正被搜到 网站做好了没人访问,这是90%企业老板最头疼的事。不是设计不够炫,也不是功能不够多,而是从底层架构到前端代码,每一步都在为SEO埋雷。我见过太多案例:花三十万做的官网,百度收录不到…

2026/9/27 0:51:05 阅读更多 →
codex-app-mirror如何15分钟发现Codex新版?“探测→比对→发布“镜像管道全拆解

codex-app-mirror如何15分钟发现Codex新版?“探测→比对→发布“镜像管道全拆解

codex-app-mirror如何15分钟发现Codex新版?"探测→比对→发布"镜像管道全拆解 【免费下载链接】codex-app-mirror 原样镜像官方 Codex 桌面应用:每 15 分钟探测、SHA256 可校验、国内直连下载、 Mac 可增量更新 | Verbatim, verifiable mirror of the off…

2026/9/27 0:50:04 阅读更多 →

日新闻

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解 【免费下载链接】spirula-studio Cross-vendor 3D Gaussian Splatting trainer - video to splat to mesh, Vulkan or CUDA. 项目地址: https://gitcode.com/GitHub_Trending/sp/spirula-studio Sp…

2026/9/27 0:00:34 阅读更多 →
SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南 模板网站太丑不够用?别急着加滤镜,那是治标不治本。很多老板盯着后台流量掉得眼红,却还在纠结首页Banner的圆角是不是3像素。这就像穿着西装去挖土,姿势不对,努力白费。我整理这份 速查手册…

2026/9/27 0:00:34 阅读更多 →
FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏 【免费下载链接】FireRed-OpenStoryline FireRed-OpenStoryline is an AI video editing agent that transforms manual editing into intention-driven directing through natural language …

2026/9/27 0:00:34 阅读更多 →

周新闻

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解 【免费下载链接】spirula-studio Cross-vendor 3D Gaussian Splatting trainer - video to splat to mesh, Vulkan or CUDA. 项目地址: https://gitcode.com/GitHub_Trending/sp/spirula-studio Sp…

2026/9/27 0:00:34 阅读更多 →
SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南 模板网站太丑不够用?别急着加滤镜,那是治标不治本。很多老板盯着后台流量掉得眼红,却还在纠结首页Banner的圆角是不是3像素。这就像穿着西装去挖土,姿势不对,努力白费。我整理这份 速查手册…

2026/9/27 0:00:34 阅读更多 →
FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏 【免费下载链接】FireRed-OpenStoryline FireRed-OpenStoryline is an AI video editing agent that transforms manual editing into intention-driven directing through natural language …

2026/9/27 0:00:34 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/25 20:29:43 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/25 20:29:31 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/26 22:52:30 阅读更多 →