
Apache Pulsar Functions 总览轻量级流式处理编程模型与实战指南【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsarPulsar Functions 是 Apache Pulsar 内置的轻量级计算框架它允许开发者直接以 Pulsar topic 为消息总线用 Java、Python 或 Go 编写函数实现消息消费、业务处理与结果发布的完整链路而无需额外部署 Storm、Flink 等外部流处理系统。本文将以 functions-overview.md 为主线结合仓库源码与官方示例系统讲解 Pulsar Functions 的编程模型、核心 API、典型用例与消息处理保证帮助读者快速掌握从编写函数到命令行部署的完整实战路径。Pulsar Functions 是什么Pulsar Functions是一类轻量级计算进程lightweight compute processes它的职责可以概括为三点从一个或多个 Pulsar topic消费消息对每条消息执行用户自定义的处理逻辑将计算结果发布到另一个 topic。这一输入 topic → 函数逻辑 → 输出 topic的模型与常见流式处理引擎一致但最大的区别在于Pulsar Functions 是 Pulsar 消息系统自身的计算基础设施运行在 Pulsar 集群内部由 Functions Worker 调度无需在集群旁边再部署一套独立的流处理系统如 Apache Storm、Apache Heron 或 Apache Flink。设计目标Pulsar Functions 的核心目标由一系列子目标支撑开发者生产力Developer productivity既可以用各语言原生的方式编写函数也可以使用 Pulsar Functions SDK 提供的接口如 Java 的FunctionI, O、Python 的Function类、Go 的pf.Function接口来获得类型支持、上下文与状态管理等高级能力易于排查问题Easy troubleshooting函数支持把日志输出到独立的 log topic并可通过bin/pulsar-admin functions系列命令查询函数状态、重启实例或获取输出让调试变得简单直接运维简洁Operational simplicity无需引入外部处理系统函数的创建、更新、删除、伸缩全部通过 Pulsar 管理命令完成。设计灵感Pulsar Functions 的形态参考了两类系统流处理引擎如 Apache Storm、Apache Heron、Apache Flink——它们提供了 topic 之间的消息处理管道Serverless 与 FaaS 云平台如 AWS Lambda、Google Cloud Functions、Azure Functions——它们提供了事件驱动、按消息执行、自动伸缩的函数即服务体验。可以这样理解Pulsar Functions 是专为以 Pulsar 作为消息总线而设计的、Lambda 风格的函数——你拿到一个回调式的process方法Pulsar 负责把输入 topic 的消息喂给它并把结果写回 Pulsar。核心编程模型Pulsar Functions 的编程模型非常简单。函数从一个或多个输入 topic接收消息每收到一条消息函数需要完成以下一类或几类任务对输入应用处理逻辑并将输出写入Pulsar 中的输出 topicoutput topicApache BookKeeper 中的状态存储state storage例如基于内置分布式 counter 做单词计数将日志写入log topic通常用于调试递增一个计数器counter。从源码层面看这一模型的载体是 Function.java 中定义的函数式接口FunctionalInterface public interface FunctionI, O { O process(I input, Context context) throws Exception; default void initialize(Context context) throws Exception {} default void close() throws Exception {} }process方法对输入 topic 的每条消息调用一次I 为输入类型O 为输出类型initialize在函数实例启动时调用一次用于初始化资源close在实例停止时调用用于释放资源。函数实例运行期间Context.java 提供了丰富的上下文能力包括getInputTopics()获取全部输入 topic 列表getOutputTopic()获取配置的输出 topicgetFunctionName() / getFunctionId() / getFunctionVersion()当前函数的名字、ID 与版本getUserConfigMap() / getUserConfigValue(key) / getUserConfigValueOrDefault(key, default)读取用户自定义配置publish(topicName, object)与newOutputMessage(topicName, schema)向指定 topic 发布消息继承自 BaseContext.java 的能力getLogger()获取日志器、incrCounter(key, amount)递增内置分布式计数器、getCounter(key)读取计数、putState / getState / deleteState操作状态存储、getSecret(name)读取密钥、recordMetric(name, value)上报自定义指标以及getTenant() / getNamespace() / getInstanceId() / getNumInstances()等运行时信息。处理链示例你可以把多个函数串接成一条处理链processing chain每个环节消费上一个环节的输出一个Python 函数监听raw-sentencestopic对输入字符串做清洗去除多余空白、全部转为小写把结果发布到sanitized-sentencestopic一个Java 函数监听sanitized-sentencestopic统计每个单词在指定时间窗口内的出现次数把结果发布到resultstopic最后一个Python 函数监听resultstopic把统计结果写入 MySQL 表。这种消息管道 独立函数的组织方式让每个处理环节都可以独立开发、独立部署、独立扩容。Word count 示例经典的 word count 是理解 Pulsar Functions 编程模型的最佳入门案例。下面这张图展示了它的工作方式sentencestopic 中的句子进入函数后每个单词对应一个 counter函数把计数结果输出到各个单词的 counter topic。Java 实现使用 Pulsar Functions SDK for Java 编写函数长这样package org.example.functions; import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.api.Function; import java.util.Arrays; public class WordCountFunction implements FunctionString, Void { // This function is invoked every time a message is published to the input topic Override public Void process(String input, Context context) throws Exception { Arrays.asList(input.split( )).forEach(word - { String counterKey word.toLowerCase(); context.incrCounter(counterKey, 1); }); return null; } }要点解读函数实现FunctionString, Void输入类型为String无输出返回null计数结果不依赖输出 topic而是写入内置状态存储每条消息被按空格拆分成单词word.toLowerCase()得到计数器 key再调用context.incrCounter(counterKey, 1)递增——这个计数器是持久化、分布式的底层由 BookKeeper 提供跨实例、跨重启保持一致仓库中 WordCountFunction.java 提供了一个更精简的变体Arrays.asList(input.split(\\s)).forEach(word - context.incrCounter(word, 1));直接用正则\s切分且没有做小写归一化。打包与部署将上述代码打包成 JAR含依赖后通过命令行部署到 Pulsar 集群$ bin/pulsar-admin functions create \ --jar target/my-jar-with-dependencies.jar \ --classname org.example.functions.WordCountFunction \ --tenant public \ --namespace default \ --name word-count \ --inputs persistent://public/default/sentences \ --output persistent://public/default/count参数说明参数含义--jar函数 JAR 文件路径需包含依赖可通过 shade 插件生成--classname实现Function接口的完整类名--tenant/--namespace函数所属租户与命名空间--name函数名称在租户/命名空间内唯一--inputs输入 topic可指定多个用逗号分隔--output输出 topic本示例函数不产生输出可留空Content-based routing 示例内容路由content-based routing是 Pulsar Functions 更进阶的典型场景函数根据消息内容把消息分发到不同的目标 topic。例如一个函数接收物品字符串作为输入根据物品类型发布到fruitstopic 或vegetablestopic如果某件物品既不是水果也不是蔬菜则向 log topic 记录一条告警日志。Python 实现用 Python 实现该路由功能如下from pulsar import Function class RoutingFunction(Function): def __init__(self): self.fruits_topic persistent://public/default/fruits self.vegetables_topic persistent://public/default/vegetables staticmethod def is_fruit(item): return item in [bapple, borange, bpear, bother fruits...] staticmethod def is_vegetable(item): return item in [bcarrot, blettuce, bradish, bother vegetables...] def process(self, item, context): if self.is_fruit(item): context.publish(self.fruits_topic, item) elif self.is_vegetable(item): context.publish(self.vegetables_topic, item) else: warning The item {0} is neither a fruit nor a vegetable.format(item) context.get_logger().warn(warning)要点解读Python SDK 中的Function基类定义在 function.py其核心是一个抽象方法process(self, input, context)——与 Java SDK 的process一一对应context.publish(topic, object)用于向任意 topic 发布消息context.get_logger()用于获取日志器并输出到 log topic输入item以字节形式如bapple出现is_fruit/is_vegetable是静态分类方法便于复用与测试。部署将代码保存为~/router.py然后通过命令行部署$ bin/pulsar-admin functions create \ --py ~/router.py \ --classname router.RoutingFunction \ --tenant public \ --namespace default \ --name route-fruit-veg \ --inputs persistent://public/default/basket-items注意Python 函数通过--py指定脚本路径--classname需要写成模块名.类名即router.RoutingFunction且不指定--output函数在process内部通过context.publish决定输出目标。函数、消息与消息类型从底层看Pulsar Functions 接收字节数组作为输入输出也是字节数组。但在支持类型化接口的语言如 Java中你可以编写类型化函数并通过以下两种方式把消息绑定到具体类型Schema Registry使用 Pulsar 内置或自定义 Schema如 JSON、Avro、Protobuf来序列化 / 反序列化消息函数声明FunctionMyType, OtherType即可直接拿到强类型对象SerDe为输入/输出类型指定序列化器类SerDeT函数框架负责在读取消息时反序列化、在发布结果时序列化。例如 Java SDK 中的 JavaSerDe.java 提供了基于 Java 序列化的默认实现而更推荐的做法是利用 Pulsar 的 Schema Registry 获得跨语言兼容的类型安全。Fully Qualified Function NameFQFN每个 Pulsar Function 都有一个完全限定函数名FQFN由三部分构成函数所属的租户tenant、命名空间namespace和函数名name。FQFN 的格式如下tenant/namespace/name例如public/default/word-count。FQFN 的存在意味着只要位于不同的命名空间你可以创建多个同名函数这为多租户场景下按命名空间隔离同名业务函数提供了支持。支持的语言目前 Pulsar Functions 支持用Java、Python、Go三种语言编写。除本文示例外各语言的详细开发指南参见 Develop Pulsar FunctionsJavaSDK 位于 pulsar-functions/api-java核心接口为FunctionI, O、WindowFunction与ContextPythonSDK 位于 pulsar-client-cpp/python/pulsar/functions核心为Function基类与Context类GoSDK 位于 pulsar-function-go/pf核心为pf.Function接口与pf.Context见 context.go。此外Pulsar Functions 还支持Window Function窗口函数形态接口定义在 WindowFunction.java用于按时间或数量窗口聚合处理。Processing guarantees处理保证Pulsar Functions 提供三种消息处理语义messaging semantics可以应用到任意函数上Delivery semantics说明At-most-once投递发给函数的每条消息可能被处理也可能不被处理至多一次消息丢失可接受、但绝不重复处理的场景使用At-least-once投递发给函数的每条消息可能被处理多次至少一次保证不丢消息但允许重复Effectively-once投递发给函数的每条消息有且仅有一个对应输出既不丢也不重有效一次这三种枚举在协议层由 Function.proto 定义其中ATLEAST_ONCE 0为默认值。为函数设置处理保证创建函数时即可指定处理保证例如下面这条pulsar-admin functions create命令创建了一个EFFECTIVELY_ONCE语义的函数$ bin/pulsar-admin functions create \ --name my-effectively-once-function \ --processing-guarantees EFFECTIVELY_ONCE \ # Other function configs--processing-guarantees的可用取值ATMOST_ONCEATLEAST_ONCEEFFECTIVELY_ONCE注意默认情况下Pulsar Functions 提供at-least-once投递保证。因此如果创建函数时未给--processing-guarantees传值函数默认使用 at-least-once 语义。更新函数的处理保证可以使用pulsar-admin functions update命令更改函数的处理保证例如$ bin/pulsar-admin functions update \ --processing-guarantees ATMOST_ONCE \ # Other function configs源码层面的语义实现不同处理保证在运行时如何落地从 PulsarSource.java 的源码可以看出函数框架在收到消息后构建Record时会根据ProcessingGuarantees决定确认ack与失败fail行为EFFECTIVELY_ONCE消息处理成功后调用consumer.acknowledgeCumulativeAsync(message)累积确认处理失败则直接抛出RuntimeException并终止从而保证输出与输入一一对应其他语义ATMOST_ONCE / ATLEAST_ONCE处理成功后调用consumer.acknowledgeAsync(message)单条确认失败时调用consumer.negativeAcknowledge(message)负确认触发重新投递配合消息去重与重试机制实现 at-least-once 乃至 effectively-once。另外FunctionConfigUtils.java 中可以看到当函数启用EFFECTIVELY_ONCE时框架会校验输入订阅类型等前置条件——例如 effectively-once 语义要求函数消费订阅满足特定约束否则会抛出异常并提示仅支持ATLEAST_ONCE的降级方案。这提醒我们在实际生产中使用 effectively-once 时需要留意订阅类型与重试配置的匹配。快速开始建议如果你希望立刻上手验证本文内容可以参考以下路径阅读 functions-quickstart.md 完成本地单机模式standalone下的第一个函数深入学习函数开发细节Schema、SerDe、状态存储、日志、用户配置见 functions-develop.md完整的命令行参考create / update / delete / get / list / status 等见 functions-cli.md 与 admin-api-functions.md生产环境部署时Functions Worker 的部署形态与 broker 共存或独立部署与配置说明见 functions-worker.md。总结Pulsar Functions 把函数即服务的理念嵌入 Pulsar 消息系统本身它没有引入外部计算框架而是以 topic 为输入/输出边界以process(input, context)为统一的编程入口配合内置的状态存储、日志 topic、分布式计数器与三种消息处理保证构成了一个开箱即用的轻量级流式处理层。无论是简单的单词统计、内容路由还是由多个函数串联的复杂处理管道你都可以在 Pulsar 生态内用 Java、Python 或 Go 独立完成这大幅降低了消息驱动的流式业务的上手与运维成本。【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考