消息队列后端流处理【免费下载链接】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),仅供参考