首页/新闻资讯/正文详情

Apache Pulsar Functions 部署实战:从本地运行到集群模式的完整指南

发布时间:2026/9/23 17:18:22 来源:云帆数科 栏目:资讯中心
Apache Pulsar Functions 部署实战:从本地运行到集群模式的完整指南
消息队列后端流处理【免费下载链接】pulsarApache 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。获取方式取决于集群的部署方式。此外如果要在部署后触发triggerPython 用户自定义函数必须在所有运行 functions worker 的机器上安装 pulsar python client。这是 Python 函数实例能够连接 broker、消费与生产消息的前提。命令行接口pulsar-admin functions 概览Pulsar Functions 的部署与管理全部通过pulsar-admin functions接口完成。该接口包含多个子命令常用的有create以集群模式部署函数trigger触发已部署的函数见下文触发 Pulsar Functionslist列出已部署的函数其他命令还包括update、delete、get、getstatus、getstats、restart、stop、start、localrun、upload、download等。从源码结构看这些子命令都在 CmdFunctions.java 中定义该类通过 JCommander 框架解析命令行参数Parameters(commandDescription Interface for managing Pulsar Functions...)声明了命令整体说明。默认参数不指定时系统如何推断管理 Pulsar Functions 时需要指定大量函数信息包括 tenant、namespace、输入/输出 topic 等。但其中部分参数在未指定时会有默认值。下表列出了完整的默认值规则参数默认值函数名Function name可对类名取任意值除 org、library 等类似类名外。例如指定--classname org.example.MyFunction时函数名为MyFunctionTenant从输入 topic 名称推导。如果输入 topic 位于marketingtenant 下即 topic 名形如persistent://marketing/{namespace}/{topicName}则 tenant 为marketingNamespace从输入 topic 名称推导。如果输入 topic 位于marketingtenant 的asianamespace 下topic 名形如persistent://marketing/asia/{topicName}则 namespace 为asia输出 topicOutput topic{输入 topic}-{函数名}-output。例如输入 topic 名为incoming、函数名为exclamation则输出 topic 名为incoming-exclamation-output订阅类型Subscription type对于at-least-once和at-most-once处理保证默认应用SHARED模式对于effectively-once保证则应用FAILOVER模式处理保证Processing guaranteesATLEAST_ONCEPulsar 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上面这条命令中函数拥有以下默认值函数名MyFunctionTenantpublicNamespacedefault订阅类型SHARED处理保证ATLEAST_ONCEPulsar service URLpulsar://localhost:6650源码视角默认参数是如何推断的默认值的推断逻辑在仓库源码中有明确实现。在 CmdFunctions.java 中NamespaceCommand.processArguments()在未指定--tenant和--namespace时分别将其置为PUBLIC_TENANT即public和DEFAULT_NAMESPACE即default。随后的FunctionCommand还支持使用--fqfnFully Qualified Function Name形如tenant/namespace/name一次性指定三者且禁止--fqfn与--tenant/--namespace/--name混用否则会抛出运行时异常。函数名的推断则由 Utils.java 中的inferMissingFunctionName完成它按.分割类名取最后一段作为函数名——例如org.example.MyFunction推断为MyFunction。若未提供 tenant/namespaceinferMissingTenant与inferMissingNamespace同样回退到public与default。此外validateFunctionConfigs 还会做完整性校验Python 与 Java 函数必须指定--classnameGo 函数不需要必须且只能指定--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 函数生效可选THREAD或PROCESS、--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并行度ParallelismPulsar 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核数KubernetesRAM字节数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参数分别被解析为Double、Long、Long类型并封装进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 服务默认不启用。enablePackagesManagementtrue packagesManagementStorageProviderorg.apache.pulsar.packages.management.storage.bookkeeper.BookKeeperPackagesStorageProvider packagesReplicas1 packagesManagementLedgerRootPath/ledgers在仓库自带的 conf/broker.conf 中可以看到这些配置的真实默认值enablePackagesManagementfalse默认关闭、packagesManagementStorageProvider默认即指向 BookKeeper 存储实现、packagesReplicas1、packagesManagementLedgerRootPath/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 FunctionsTrigger如果一个 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。/DSMLparameter /DSMLinvoke /DSMLtool_calls赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Functions 快速上手指南从本地运行到集群部署的完整实践Apache Pulsar Functions 快速上手指南从本地运行到集群部署的完整实践 本篇技术指南以 Apache Pulsar 官方入门文档为基础带消息队列后端流处理Google IMA SDK WebHTML5客户端广告插入完整集成指南Google IMA SDK WebHTML5客户端广告插入完整集成指南 本指南基于 ima sdk web guide.md https://link.g消息队列后端流处理Apache Pulsar 裸机多集群部署完整指南从 ZooKeeper、BookKeeper 到 Broker 的实战部署Apache Pulsar 裸机多集群部署完整指南从 ZooKeeper、BookKeeper 到 Broker 的实战部署 导读 本文是基于当前 Apach消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

3个源码图解原理带你搞定繁体版从入门到实战
3个源码图解原理带你搞定繁体版从入门到实战

3个源码图解原理带你搞定繁体版从入门到实战 学会语法却不知怎么搭项目,这是无数开发者卡在“半吊子”阶段的死穴。你背熟了 API,却连一个完整的繁体转换模块都写不出来。今天不讲虚的,直接扒开【繁体版】转换的核心源码,用图解原理拆解底层逻辑,让… · 2026/9/23 17:18:16

《HarmonyOS 7 应用上架与隐私合规工程化》01:权限声明、运行时申请与 Release Profile 为什么总对不上【鸿蒙心迹】
《HarmonyOS 7 应用上架与隐私合规工程化》01:权限声明、运行时申请与 Release Profile 为什么总对不上【鸿蒙心迹】

开发机上一切正常,提交审核直接被打回来:权限声明和 Release Profile 对不上。做了个生活记录 App,开发机上跑了三个月,功能都正常。提交应用市场审核,等了两天,收到驳回通知:“敏感隐私权限未提… · 2026/9/23 17:18:09

3分钟搞定Excel数据透视手写实现
3分钟搞定Excel数据透视手写实现

3分钟搞定Excel数据透视手写实现 官方文档翻了三遍还是晕?别慌。Excel数据透视表看着复杂,其实底层逻辑就三步:聚合、分组、求和。今天不聊虚的,咱们直接上手,用Python代码把这套逻辑跑通。哪怕你是刚入行的房建工程师,或者对机器学习… · 2026/9/23 17:18:09

Python多线程穷举ZIP/RAR/7Z密码工具实战指南
Python多线程穷举ZIP/RAR/7Z密码工具实战指南

简介:这是一套基于Python实现的多线程可视化压缩包密码破解工具,面向信息安全初学者、CTF备赛者及渗透测试爱好者,用于学习密码学基础、暴力破解原理与多线程编程实践。资源包含238个文件,主体为9个核心Python脚本(含G… · 2026/9/23 18:33:22

PP-LCNet 图像分类实战指南:基于 PaddleHub 使用 pplcnet_x2_5_imagenet 完成推理与服务部署
PP-LCNet 图像分类实战指南:基于 PaddleHub 使用 pplcnet_x2_5_imagenet 完成推理与服务部署

PP-LCNet 图像分类实战指南:基于 PaddleHub 使用 pplcnet_x2_5_imagenet 完成推理与服务部署 【免费下载链接】PaddleFormers PaddleFormers is an easy-to-use library of pre-trained large language model zoo based on PaddlePaddle. 项目地址: https://gitco… · 2026/9/23 18:33:16

面试被问原理答不上?一文搞懂免费酒店管理系统
面试被问原理答不上?一文搞懂免费酒店管理系统

面试被问原理答不上?一文搞懂免费酒店管理系统 面试时,面试官轻飘飘问一句:“讲下你做的酒店管理系统,核心逻辑怎么流转?”结果你卡壳了。脑子一片空白,只记得写了增删改查,却说不清库存扣减、房态同步、并发锁死这些底层原理。… · 2026/9/23 18:33:10

OpenJarvis Skills系统完全指南:13000+社区技能如何教会AI用工具
OpenJarvis Skills系统完全指南:13000+社区技能如何教会AI用工具

OpenJarvis Skills系统完全指南:13000社区技能如何教会AI用工具 【免费下载链接】OpenJarvis Personal AI, On Personal Devices 项目地址: https://gitcode.com/gh_mirrors/op/OpenJarvis OpenJarvis 是一个运行在个人设备上的开源个人 AI 智能体框架&#… · 2026/9/23 18:33:09

种植牙医院排名系统卡顿?3招性能优化让查询秒出
种植牙医院排名系统卡顿?3招性能优化让查询秒出

种植牙医院排名系统卡顿?3招性能优化让查询秒出 刚接手一个医疗垂直搜索项目,核心需求是展示【种植牙医院排名】。上线第一天就炸了,后台日志全是超时报警。用户反馈说,搜索“北京朝阳区种植牙哪家好”时,页面加载要等8秒,转圈圈转到怀疑人生。我盯着… · 2026/9/23 18:33:03

Somin配置卡死救急:3个实战项目避坑指南
Somin配置卡死救急:3个实战项目避坑指南

Somin配置卡死救急:3个实战项目避坑指南 刚接触Somin的朋友,大概率经历过这种绝望:明明照着教程敲命令,环境就是起不来,报错信息像天书一样滚过去,卡在那儿半天动不了。这种“配置环境就卡半天”的体验,直接劝退了一半想入坑的人。… · 2026/9/23 18:33:03

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码