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

Apache Pulsar Functions Worker 部署与管理实战:与 Broker 合跑与独立运行两种模式详解

发布时间:2026/9/23 2:49:08 来源:云帆数科 栏目:资讯中心
Apache Pulsar Functions Worker 部署与管理实战:与 Broker 合跑与独立运行两种模式详解
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文以 Apache Pulsar 的functions-worker为核心系统讲解 Pulsar Functions 集群模式运行的基础组件——Functions Worker 的两种部署形态与 Broker 合跑、独立运行涵盖配置要点、安全加固、代理转发与故障排查并结合 functions-worker.md 与仓库源码、配置文件给出可落地的实操指引。读完本文你将掌握如何根据资源隔离、Kubernetes 与多集群诉求选择合适的运行模式并正确配置functions_worker.yml、broker.conf与proxy.conf完成部署。说明文中--- Service Urls ---线条在架构图中代表 Pulsar 客户端与管理端用于连接 Pulsar 集群的服务 URL。一、Functions Worker 是什么Pulsarfunctions-worker是负责在集群模式下运行 Pulsar Functions 的逻辑组件负责函数包的存储与分发、函数元数据管理assignment、调度等、运行时Runtime的创建与监控以及 Functions/Source/Sink 相关的 Admin REST 接口。要使用 Pulsar Functions你需要先掌握如何搭建 functions-worker并配置 Functions runtime。在集群中你可以根据需求选择两种运行方式与 Broker 合跑Run with brokers独立运行Run separately从源码结构看functions-worker的核心配置模型位于 WorkerConfig.java该文件定义了本文涉及的全部配置参数及其默认值默认配置文件为 conf/functions_worker.yml。二、与 Broker 合跑Run Functions-worker with brokers2.1 启用合跑在conf/broker.conf中设置functionsWorkerEnabledtruebroker 启动时即会内嵌启动 functions-workerfunctionsWorkerEnabledtrue对应配置项在 ServiceConfiguration.java 中定义默认值为false。一旦开启你需要继续配置 conf/functions_worker.yml 来定制 functions-worker 的行为。合跑模式下大部分设置已从 broker 配置继承如 configurationStore、认证设置等但有两个必选设置需要格外注意配置项含义默认值建议numFunctionPackageReplicas函数包的副本数1适合 standalone生产环境为保证高可用建议设为 ≥ 2initializedDlogMetadata是否在运行时初始化分布式日志元数据false若设为true须先用bin/pulsar initialize-cluster-metadata命令完成初始化2.2 配置 BookKeeper 认证如果 BookKeeper 集群开启了认证需要在functions_worker.yml中配置以下三项bookkeeperClientAuthenticationPluginBookKeeper 客户端认证插件名bookkeeperClientAuthenticationParametersNameBookKeeper 客户端认证插件参数名bookkeeperClientAuthenticationParametersBookKeeper 客户端认证插件参数2.3 配置 Stateful-Functions状态函数若需要使用putState()、queryState()等状态相关接口按以下步骤启用第 1 步在 BookKeeper 中启用 streamStorage 服务该服务当前基于 NAR 包在 conf/bookkeeper.conf 中配置extraServerComponentsorg.apache.bookkeeper.stream.server.StreamStorageLifecycleComponentbookie 启动后可用telnet验证服务是否正常telnet localhost 4181成功输出Trying 127.0.0.1... Connected to localhost. Escape character is ^].第 2 步在functions_worker.yml中开启状态存储stateStorageServiceUrl: bk://bk-service-url:4181其中bk-service-url指向 BookKeeper table service 的服务 URL。对应源码字段见 WorkerConfig.java默认状态存储实现类为BKStateStoreProviderImpl。2.4 启动并验证配置完成后启动或重启 broker。随后用curl验证 functions-worker 是否正常运行curl broker-ip:8080/admin/v2/worker/cluster返回集群中活跃 function worker 的列表例如[{workerId:worker-id,workerHostname:worker-hostname,port:8080}]三、独立运行 Functions WorkerRun separately重要独立运行模式下务必保持functionsWorkerEnabledfalse避免误与 broker 合跑。同时使用pulsar-admin或客户端管理函数时需根据 Worker parameters 中设置的workerHostname与workerPort生成--admin-url。3.1 Worker 参数配置项类型/含义说明workerIdstring跨集群唯一用于标识一台 worker 机器未设置时源码会根据workerHostname-workerPort自动生成workerHostnameworker 机器的主机名未设置时源码会通过本机 DNS 解析获取workerPortworker server 监听端口默认6750workerPortTlsworker server 的 TLS 监听端口默认6751对应源码见 WorkerConfig.java。3.2 函数包参数numFunctionPackageReplicas函数包副本数默认1。3.3 函数元数据参数配置项含义pulsarServiceUrlbroker 集群的 Pulsar 服务 URL二进制协议pulsarWebServiceUrlbroker 集群的 Pulsar Web 服务 URLHTTP 协议pulsarFunctionsCluster设为你的 Pulsar 集群名与 broker 配置中的clusterName一致若 broker 集群开启了认证还应配置 worker 与 broker 通信的认证信息clientAuthenticationPluginclientAuthenticationParameters注意合跑模式下 worker 可直接继承 broker 的认证配置独立运行模式则必须显式配置。3.4 定制 Java 运行时参数通过additionalJavaRuntimeArguments可向每个由 functions worker 启动的进程追加 JVM 命令行参数additionalJavaRuntimeArguments: [-XX:ExitOnOutOfMemoryError,-Dfoobar]典型用途添加 JVM flags如-XX:ExitOnOutOfMemoryError传入自定义系统属性如-Dlog4j2.formatMsgNoLookups注意该特性仅适用于 Process 与 Kubernetes 运行时。对应源码字段见 WorkerConfig.java。3.5 安全设置如需启用 Functions Worker 的安全能力应按顺序配置启用 TLS 传输加密启用认证提供方启用授权提供方启用端到端加密3.5.1 启用 TLS 传输加密useTLS: true pulsarServiceUrl: pulsarssl://localhost:6651/ pulsarWebServiceUrl: https://localhost:8443 tlsEnabled: true tlsCertificateFilePath: /path/to/functions-worker.cert.pem tlsKeyFilePath: /path/to/functions-worker.key-pk8.pem tlsTrustCertsFilePath: /path/to/ca.cert.pem # Pulsar 客户端用于向 broker 认证时使用的受信证书路径 brokerClientTrustCertsFilePath: /path/to/ca.cert.pem详细内容参见 Transport Encryption using TLS。3.5.2 启用认证提供方在functions_worker.yml中设置authenticationEnabled: true authenticationProviders: [ provider1, provider2 ]注意将providers list替换为你实际启用的提供方。TLS 认证详见 TLS AuthenticationbrokerClientAuthenticationPlugin: org.apache.pulsar.client.impl.auth.AuthenticationTls brokerClientAuthenticationParameters: tlsCertFile:/path/to/admin.cert.pem,tlsKeyFile:/path/to/admin.key-pk8.pem authenticationEnabled: true authenticationProviders: [org.apache.pulsar.broker.authentication.AuthenticationProviderTls]SASL 认证如需可在properties下配置saslJaasClientAllowedIds与saslJaasServerSectionNameproperties: saslJaasClientAllowedIds: .*pulsar.* saslJaasServerSectionName: BrokerToken 认证详见 Token Authenticationproperties: tokenSecretKey: file://my/secret.key # 若使用公私钥 # tokenPublicKey: file:///path/to/public.key注意密钥文件必须是 DER 编码。3.5.3 启用授权提供方启用授权需要配置authorizationEnabled、authorizationProvider与configurationMetadataStoreUrl。认证提供方通过连接configurationMetadataStoreUrl获取命名空间策略authorizationEnabled: true authorizationProvider: org.apache.pulsar.broker.authorization.PulsarAuthorizationProvider configurationMetadataStoreUrl: meta-type:configuration-metadata-store-url还需配置超级用户角色列表可访问任意 Admin APIsuperUserRoles: - role1 - role2 - role33.5.4 启用端到端加密端到端加密使用应用配置的公私钥对进行加解密只有持有有效密钥的消费者才能解密消息。在命令行中通过--producer-config指定详见 security-encryption.md。CryptoConfig的相关字段已并入ProducerConfig可配置字段如下public class CryptoConfig { private String cryptoKeyReaderClassName; private MapString, Object cryptoKeyReaderConfig; private String[] encryptionKeys; private ProducerCryptoFailureAction producerCryptoFailureAction; private ConsumerCryptoFailureAction consumerCryptoFailureAction; }字段取值producerCryptoFailureAction生产端加密失败时的动作FAIL、SENDconsumerCryptoFailureAction消费端解密失败时的动作FAIL、DISCARD、CONSUME3.5.5 BookKeeper 认证若 BookKeeper 集群启用了认证配置bookkeeperClientAuthenticationPluginBookKeeper 客户端认证插件名bookkeeperClientAuthenticationParametersNameBookKeeper 客户端认证插件参数名bookkeeperClientAuthenticationParametersBookKeeper 客户端认证插件参数3.6 启动 Functions Worker后台启动配合 pulsar-daemon 与nohupbin/pulsar-daemon start functions-worker前台启动bin/pulsar functions-worker3.7 为 Functions Worker 配置代理Proxy独立运行模式下Admin REST 端点被拆分成两套集群functions、function-worker、source、sink端点由functions-worker集群提供服务其余端点仍由 broker 集群提供服务。因此你需要让pulsar-admin使用正确的服务 URL。为统一入口可以启动一个代理集群来按需路由 Admin REST 请求。若尚未搭建代理集群可参考 administration-proxy.md 完成搭建。在 conf/proxy.conf 中配置以下两项以启用路由functionWorkerWebServiceURLpulsar-functions-worker-web-service-url functionWorkerWebServiceURLTLSpulsar-functions-worker-web-service-url四、两种模式对比与选型建议与 Broker 合跑更便捷但独立运行能为Process或Thread模式运行的函数提供更好的资源隔离。选择建议如下。推荐使用Run-with-Broker模式a) 在Process或Thread模式下运行函数时不需要资源隔离b) 配置 functions-worker 在 Kubernetes 上运行函数资源隔离问题由 Kubernetes 解决。推荐使用Run-separately模式a) 没有 Kubernetes 集群b) 希望将函数与 broker 分开运行。五、常见问题排查Troubleshooting报错信息Namespace missing local cluster name in clusters listFailed to get partitioned topic metadata: org.apache.pulsar.client.api.PulsarClientException$BrokerMetadataException: Namespace missing local cluster name in clusters list: local_clusterxyz nspublic/functions clusters[standalone]出现该错误通常是以下两种情况之一a) broker 以functionsWorkerEnabledtrue启动但conf/functions_worker.yaml中的pulsarFunctionsCluster未设置为正确的集群b) 搭建跨地域复制geo-replicationPulsar 集群时开启functionsWorkerEnabledtrue一个集群的 broker 正常另一集群的 broker 异常。解决方法设置functionsWorkerEnabledfalse禁用 Functions Worker并重启 broker。查询public/functions命名空间的当前集群列表bin/pulsar-admin namespaces get-clusters public/functions检查集群是否在列表中若不在将其加入并更新bin/pulsar-admin namespaces set-clusters --clusters existing-clusters,new-cluster public/functions集群设置成功后将functionsWorkerEnabledtrue重新启用 Functions Worker。在conf/functions_worker.yml中设置正确的pulsarFunctionsCluster并重启 broker。附functions_worker.yml关键配置速览以仓库自带 conf/functions_worker.yml 为例常用配置项整理如下未列出的配置项均保留默认值即可配置项示例值说明workerIdstandaloneworker 标识跨集群唯一workerHostnamelocalhostworker 主机名workerPort/workerPortTls6750/6751worker HTTP/HTTPS 监听端口configurationMetadataStoreUrlzk:localhost:2181配置元数据存储地址numFunctionPackageReplicas1函数包副本数pulsarServiceUrl/pulsarWebServiceUrlpulsar://localhost:6650/http://localhost:8080元数据管理客户端连接的 Pulsar 服务地址pulsarFunctionsNamespacepublic/functions函数元数据 topic 所在命名空间pulsarFunctionsClusterstandalone集群名须与 brokerclusterName一致schedulerClassNameorg.apache.pulsar.functions.worker.scheduler.RoundRobinScheduler函数调度器实现functionRuntimeFactoryClassNameorg.apache.pulsar.functions.runtime.process.ProcessRuntimeFactory运行时工厂Process/Thread/KubernetesauthenticationEnabled/authorizationEnabledfalse是否启用认证/授权tlsEnabledfalse是否启用 worker TLSuseTlsfalseworker 内嵌 Pulsar 客户端是否使用 TLS 连接 brokerstateStorageServiceUrlbk://localhost:4181状态存储服务 URL默认注释按需开启initializedDlogMetadatafalse是否由运行时初始化分布式日志元数据connectorsDirectory/functionsDirectory./connectors/./functions内置 Connector/Function 目录所有上述字段均可在 WorkerConfig.java 中逐一找到对应定义与默认值自定义参数可通过properties映射注入如 Token 认证所需的tokenPublicKey等。结合实际部署形态合跑或独立与安全需求按本文各小节配置即可完成 Functions Worker 的部署与管理。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Functions Worker 部署与运维完整指南随 Broker 共置与独立部署双模式实战Apache Pulsar Functions Worker 部署与运维完整指南随 Broker 共置与独立部署双模式实战 Apache Pulsar 的 F消息队列后端流处理Apache Flink CDC 独立部署模式详解与实践指南Apache Flink CDC 独立部署模式详解与实践指南 概述 Apache Flink CDC 是基于 Apache Flink 构建的变更数据捕获框架后端数据集成大数据流处理变更数据捕获数据同步Apache Flink CDC 独立部署模式详解Apache Flink CDC 独立部署模式详解 概述 Apache Flink CDC 是基于 Flink 构建的变更数据捕获 Change Data Ca后端数据集成大数据流处理变更数据捕获数据同步创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

vGPU与GPU直通本质区别:从硬件隔离到AI推理选型指南
vGPU与GPU直通本质区别:从硬件隔离到AI推理选型指南

1. 项目概述:为什么今天必须重新理解 GPU 虚拟化的底层逻辑你是不是也遇到过这样的场景:在 VMware Workstation 里装了个 Windows 11 虚拟机,想跑个 Stable Diffusion WebUI,结果点开“显示设置”发现显卡选项灰掉——提示“在此主… · 2026/9/23 2:49:01

Click 参数详解:用 @click.option 与 @click.argument 为命令行命令注入输入
Click 参数详解:用 @click.option 与 @click.argument 为命令行命令注入输入

Click 参数详解:用 click.option 与 click.argument 为命令行命令注入输入 【免费下载链接】Tutorial-Codebase-Knowledge Pocket Flow: Codebase to Tutorial 项目地址: https://gitcode.com/gh_mirrors/tu/Tutorial-Codebase-Knowledge 本篇技术指南深入讲… · 2026/9/23 2:49:01

类型安全容器设计:从泛型约束到工程落地实践
类型安全容器设计:从泛型约束到工程落地实践

在写代码这些年里,我越来越觉得“类型安全”这四个字被低估了。很多人对类型安全的认知停留在“编译不过就改一改”的层面,但真正把它落实到容器设计上,能少踩的坑远比想象中多。今天想用自己的实际项目经验,聊聊类型安全容器设计… · 2026/9/23 2:49:01

智能化系统集成项目经理培训机构推荐:从报名学习到考试拿证,报考全攻略
智能化系统集成项目经理培训机构推荐:从报名学习到考试拿证,报考全攻略

智慧楼宇、智慧园区项目的落地,离不开系统集成的整体统筹。智能化系统集成项目经理作为智能化项目的”总设计师总调度”,是行业中的高端人才。本文给你一份完整的智能化系统集成项目经理报考全攻略。 一、智能化系统集成项目经理是做什么的? … · 2026/9/23 3:36:18

Apache Druid Coordinator 节点配置完全指南:运行参数、动态配置与源码级原理
Apache Druid Coordinator 节点配置完全指南:运行参数、动态配置与源码级原理

数据库数据分析OLAP大数据实时分析数据仓库后端 【免费下载链接】druid Apache Druid: a high performance real-time analytics database. 项目地址: https://gitcode.com/gh_mirrors/druid7/druid 点击查看 免费下载 本篇技术指南围绕 Apache Druid 集群中的 Coo… · 2026/9/23 3:36:18

搞定贝努鸟:3步重构解决版本升级后API全变痛点
搞定贝努鸟:3步重构解决版本升级后API全变痛点

搞定贝努鸟:3步重构解决版本升级后API全变痛点 上周刚把项目里的核心模块从 v2 升级到 v3,结果一跑测试,满屏红叉。最让人头大的是,原本封装好的 BirdEngine 接口在 v3 里直接重构了, fetch() 变成了… · 2026/9/23 3:36:12

数据库工程师培训机构推荐:从报名学习到考试拿证,报考全攻略
数据库工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

数据是企业的核心资产,数据库工程师是管理和守护这些资产的”数据管家”。从银行交易到电商订单,数据库工程师是IT系统不可或缺的技术岗位。本文给你一份完整的数据库工程师报考全攻略。 一、数据库工程师是做什么的? 数据库工程师是负责数据… · 2026/9/23 3:36:12

网络安全架构师培训机构推荐:从报名学习到考试拿证,报考全攻略
网络安全架构师培训机构推荐:从报名学习到考试拿证,报考全攻略

在企业安全体系从”单点防护”走向”整体防御”的今天,网络安全架构师作为安全体系的顶层设计者,是行业中的高端稀缺人才。本文给你一份完整的网络安全架构师报考全攻略。 一、网络安全架构师是做什么的? 网络安全架构师是负责企业网络安全体… · 2026/9/23 3:36:12

高级智能家居系统工程师培训机构推荐:从报名学习到考试拿证,报考全攻略
高级智能家居系统工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

智能家居行业正从”单品智能”走向”全屋智能”,高级智能家居系统工程师作为方案设计与项目交付的骨干,价值愈发凸显。本文给你一份完整的高级智能家居系统工程师报考全攻略。 一、高级智能家居系统工程师是做什么的? 高级智能家居系统工程师… · 2026/9/23 3:36:12

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

了解更多?预约专属演示

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

企业微信二维码