消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar Functions 是 Pulsar 提供的轻量级、Lambda 风格的计算能力允许用户以函数方式消费并处理消息。本文围绕 Pulsar 2.2.1 版本官方文档《Deploying and managing Pulsar Functions》展开系统讲解本地运行local run与集群模式cluster mode两种部署方式、FQFN 命名规则、默认参数、parallelism 与实例资源分配以及使用pulsar-admin命令行完成函数创建、更新、触发和状态管理的完整流程。读完本文你将能够独立把一个 Java/Python/Go 函数部署到本地或 Pulsar 集群并掌握触发验证、并行度调整等日常运维技能。两种部署模式总览Pulsar Functions 当前支持两种部署模式区别在于函数代码实际运行的位置模式说明Local run 模式本地运行函数运行在你本地环境中例如你的笔记本或某台 EC2 实例上Cluster 模式集群运行函数运行在Pulsar 集群内部与 Pulsar broker 运行在同一批机器上从源码结构看Pulsar Functions 在设计上考虑了扩展性。CLI 层由 CmdFunctions.java 统一承载其中注册了localrun、create、delete、update、get、list、trigger、querystate、putstate等子命令未来若社区新增部署选项可在该命令框架上继续扩展。如需为项目贡献新的部署方式建议先与 Pulsar 开发者社区联系。部署前的要求要部署和管理 Pulsar Functions你需要先有一个正在运行的 Pulsar 集群可选方式包括在本机运行 standalone 集群在 Kubernetes、Amazon Web Services、裸机bare metal、DC/OS 等环境部署集群。如果你运行的不是 standalone 集群需要先获取集群的 service URL获取方式取决于集群的部署方式。此外若要部署和触发 Python 用户自定义函数需要先安装 Pulsar Python client。命令行接口与 FQFNPulsar Functions 的部署与管理统一通过pulsar-admin functions接口完成其中包含create集群模式部署函数、trigger触发函数、list列出已部署函数等命令。Fully Qualified Function NameFQFN每个 Pulsar Function 都有一个Fully Qualified Function NameFQFN由三个元素组成函数所属的 tenant、namespace 以及函数名格式如下tenant/namespace/nameFQFN 让你可以在不同的 namespace 中创建多个同名函数。例如marketing/asia/MyFunction与finance/emea/MyFunction互不冲突。在 CLI 层面CmdFunctions.java 中的parseFullyQualifiedFunctionName方法会按/拆分 FQFN 并校验必须恰好为 3 段否则抛出ParameterException同时要求--fqfn不能与--tenant/--namespace/--name混用String[] args fqfn.split(/); if (args.length ! 3) { throw new ParameterException(Fully qualified function names (FQFNs) must be of the form tenant/namespace/name); } else { functionConfig.setTenant(args[0]); functionConfig.setNamespace(args[1]); functionConfig.setName(args[2]); }默认参数管理 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 为asiaOutput topic{输入 topic}-{函数名}-output。例如输入 topic 为incoming、函数名为exclamation输出 topic 为incoming-exclamation-outputSubscription type对于 at-least-once 和 at-most-once 的 processing guarantees默认应用SHARED对于 effectively-once默认应用FAILOVERProcessing guaranteesATLEAST_ONCEPulsar service URLpulsar://localhost:6650这些推断逻辑在 Utils.java 中有明确实现inferMissingFunctionName取类名的最后一个.分段作为函数名inferMissingTenant与inferMissingNamespace分别回落到publictenant 与defaultnamespacepublic static void inferMissingFunctionName(FunctionConfig functionConfig) { String[] domains functionConfig.getClassName().split(\\.); if (domains.length 0) { functionConfig.setName(functionConfig.getClassName()); } else { functionConfig.setName(domains[domains.length - 1]); } } public static void inferMissingTenant(FunctionConfig functionConfig) { functionConfig.setTenant(PUBLIC_TENANT); } public static void inferMissingNamespace(FunctionConfig functionConfig) { functionConfig.setNamespace(DEFAULT_NAMESPACE); }其中PUBLIC_TENANT与DEFAULT_NAMESPACE即public与default。默认参数的使用示例以下create命令$ bin/pulsar-admin functions create \ --jar my-pulsar-functions.jar \ --classname org.example.MyFunction \ --inputs my-function-input-topic1,my-function-input-topic2创建出的函数会为以下参数自动填充默认值函数名MyFunction、tenantpublic、namespacedefault、subscription typeSHARED、processing guaranteesATLEAST_ONCE、Pulsar service URLpulsar://localhost:6650。注意CLI 在真正提交前会通过validateFunctionConfigs校验配置例如 Java/Python 函数必须提供--classname且--jar、--py、--go三者必须且只能指定一个CmdFunctions.java。同时--jar/--py/--go还支持http、file、function包管理服务 URL等远程包地址由Utils.isFunctionPackageUrlSupported判断。Local run 模式在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 集群。如果希望在本地运行函数但连接非本地的 Pulsar 集群可以用--broker-service-url指定 broker URL$ bin/pulsar-admin functions localrun \ --broker-service-url pulsar://my-cluster-host:6650 \ # 其他函数参数从实现看LocalRunnerCmdFunctions.java会将序列化后的FunctionConfig连同--broker-service-url、--use-tls、--client-auth-plugin、--client-auth-params、--secrets-provider-classname、--metrics-port-start、--runtimeTHREAD或PROCESS等参数透传给$PULSAR_HOME/bin/function-localrunner脚本启动本地实例并通过ProcessBuilder阻塞等待进程结束。也就是说local run 是在本机拉起一个独立的函数运行时进程与集群中的函数运行方式共用同一套运行时配置模型。Cluster 模式在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-1CreateFunction在执行时会判断包地址类型若为远程 URL 则调用createFunctionWithUrl否则调用createFunction(functionConfig, userCodeFile)将本地文件上传到集群。更新集群模式函数使用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-topicUpdateFunction同样支持 URL 包与本地文件两种方式并可通过--update-auth-data控制是否更新认证数据。更新成功后会输出Updated successfully。并行度ParallelismPulsar Functions 以称为instance实例的进程运行。默认情况下一个函数只运行单个实例在 local run 模式下只能运行单实例。你可以在创建函数时通过--parallelism标志指定并行度即要运行的实例数量$ bin/pulsar-admin functions create \ --parallelism 3 \ # 其他函数信息对于已创建的函数可通过update接口调整并行度$ bin/pulsar-admin functions update \ --parallelism 5 \ # 其他函数信息如果通过 YAML 指定函数配置使用parallelism参数。示例配置文件# function-config.yaml parallelism: 3 inputs: - persistent://public/default/input-1 output: persistent://public/default/output-1 # 其他参数对应的更新命令$ bin/pulsar-admin functions update \ --function-config-file function-config.yaml在 CLI 实现中--function-config-file会通过CmdUtils.loadConfig(fnConfigFile, FunctionConfig.class)把 YAML 反序列化为FunctionConfig因此 YAML 中可配置的字段与命令行参数一一对应CmdFunctions.java。在集群模式下并行度最终会换算为函数实例数量由 FunctionRuntimeManager.java 等 worker 侧组件调度。函数实例资源分配在集群模式下你可以为每个函数 实例 指定资源资源指定方式运行时CPU核心数Docker即将支持RAM字节数Process、DockerDisk space字节数Docker下面是一个为函数分配 8 核、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 内存。做资源规划时务必把并行度即实例数计入总量。从 Resources.java 的实现看未显式指定资源时存在内置默认值CPU 默认 1 核、RAM 默认 1 GB1073741824L字节、磁盘默认 10 GB10737418240L字节mergeWithDefault方法会对空缺字段逐一回填默认值。CLI 端--cpu、--ram、--disk参数会写入Resources对象并挂到FunctionConfig.resources上。注意--cpu与--disk目前主要适用于 Docker runtime--ram适用于 Process/Docker runtimeCLI 参数描述中亦有注明。触发 Pulsar FunctionsTrigger如果 Pulsar Function 运行在 集群模式下你可以随时通过命令行触发trigger它。触发函数意味着你向函数发送一个带有特定值的消息并通过命令行获得函数的输出如果有的话。触发一个函数本质上与向函数的某个输入 topic 生产消息来调用它没有区别。pulsar-admin functions trigger命令是一个便捷机制让你无需借助pulsar-client工具或某种语言的原生 client library就能向函数发送消息。端到端触发示例首先从一个简单的 Python 函数 开始它根据输入返回一段字符串# myfunc.py def process(input): return This function has been triggered with a value of {0}.format(input)在 local run 模式下运行该函数$ 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 # 无限监听现在触发该函数$ 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。从实现看TriggerFunctionCmdFunctions.java只依赖 FQFN 三要素定位函数--trigger-value直接作为消息内容注入它还支持--trigger-file从文件读取触发数据以及可选的--topic指定要注入数据的消费 topic。runCmd会调用getAdmin().functions().triggerFunction(...)并把函数返回值打印到标准输出。更多管理命令除创建、更新、触发外pulsar-admin functions还提供以下日常管理命令CmdFunctions.java 中均有对应实现list列出指定 tenant/namespace 下运行的所有函数--tenant/--namespace缺省回落public/defaultget获取函数的FunctionConfig详情JSON 格式化输出status/getstatus查看函数可按--instance-id查单个实例的运行状态stats获取函数或指定实例的运行指标restart/stop/start重启、停止、启动函数实例未指定--instance-id时作用于全部实例delete删除集群中运行的一个函数querystate/putstate读取/写入与函数关联的 key-value 状态querystate支持-k/--key指定键名、-w/--watch持续监听putstate通过-s/--state传入 JSON 状态对象upload/download向 Pulsar 上传/从 Pulsar 下载函数包文件。这些命令统一以pulsar-admin functions subcommand的形式使用例如$ bin/pulsar-admin functions list --tenant public --namespace default $ bin/pulsar-admin functions get --tenant public --namespace default --name myfunc $ bin/pulsar-admin functions delete --tenant public --namespace default --name myfunc总结Pulsar Functions 的部署与运维整体围绕pulsar-admin functions命令展开local run 模式适合在开发机上快速验证函数逻辑函数以独立进程运行在本机cluster 模式则将函数代码上传到集群与 broker 并排运行适合生产环境。理解 FQFNtenant/namespace/name、默认参数推导、--parallelism并行度与按实例计算的资源配额CPU/RAM/Disk是正确运维函数的关键trigger命令则为验证函数行为提供了不依赖客户端工具的直接手段。配合 reference-pulsar-admin 与 functions-api 文档你可以进一步了解每个参数的完整语义与函数编写 API。赞分享消息队列后端流处理【免费下载链接】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 官方文档 functions d消息队列后端流处理Apache Pulsar Functions 部署与管理实战本地运行、集群模式、并行度与触发机制Apache Pulsar Functions 部署与管理实战本地运行、集群模式、并行度与触发机制 导读 Pulsar Functions 是 Apache消息队列后端流处理上一篇告别手动打包Vue-Web-Extension生产环境自动化工作流全攻略下一篇代码重构实战指南从能跑到优雅的蜕变之路创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
