6.3 Functions、连接器、事务与演进


6.3 Functions、连接器、事务与演进

本节摘要:远洋水域有三样装备:Functions 把消费逻辑简化成一段部署进集群的代码,连接器把外部系统接进水道,事务让跨主题操作原子化;再把视线放远,Kubernetes 化部署与存算进一步解耦是这条水道明确的下一程。本节是生态全景加演进路标,帮你在功能清单与团队承接力之间做出清醒取舍。

漂流终点的加工厂:Functions 的定位

Pulsar Functions 是内建在集群里的轻量计算:你只写一段处理函数——进来一条消息、出去一条消息(或零条)——部署、扩缩容、失败重试全由集群代管。它不是 Flink 那样的流处理引擎,没有窗口聚联合并之类的重装备;它的定位是"把胶水代码搬进水道":订单事件的格式转换、过滤、简单富化,这类逻辑以前要单独起一个消费者服务,现在一段函数搞定。

一段 Python 函数把订单事件过滤并富化后写入新主题:

from pulsar import Function class EnrichOrder(Function): def process(self, input, context): evt = json.loads(input) # 过滤:测试订单不入分析流 if evt.get("flag") == "TEST": return None # 富化:补金额档位标签 amount = evt.get("amount", 0) evt["tier"] = "HIGH" if amount >= 10000 else "NORMAL" context.publish("persistent://trade-order/analytics/order-enriched", json.dumps(evt)) return None

部署与验证一条命令完成,函数的扩缩容就是改实例数——因为它本质是集群托管的一组消费者:

# 部署函数:订阅订单主题,Python 运行时,两实例并行 $ pulsar-admin functions create \ --tenant trade-order --namespace transaction \ --name enrich-order --classname EnrichOrder \ --py order_enrich.py --inputs persistent://trade-order/transaction/order-events \ --output persistent://trade-order/analytics/order-enriched \ --parallelism 2 # 看运行状态与吞吐 $ pulsar-admin functions status --tenant trade-order \ --namespace transaction --name enrich-order

什么时候用 Functions、什么时候老实用消费者服务?判断依据是状态与复杂度:无状态或轻状态的转换逻辑用函数,省一整套服务运维;需要窗口聚合、多流关联、恰好一次语义的复杂管道,交给专业流处理引擎(Flink 直连 Pulsar 是社区成熟玩法)。把 Functions 当轻量胶水,别把它当全家桶。

一、连接器与 Pulsar IO:水道与外部世界的桥

连接器(Pulsar IO)分两个方向:Source 把外部数据拉进主题(数据库变更捕获、文件、MQTT 设备上行),Sink 把主题数据推去外部(搜索引擎、数据湖、缓存)。对订单场景,典型装配是把订单库的变更流经 Source 引进主题、把富化后的事件经 Sink 落进检索引擎——两端的代码都是现成连接器加一份配置,不再手写搬运服务。

连接器本身也是函数的变体(Pipeline 的一种),运行在同一个托管体系里。选型提醒:连接器适合"标准协议加简单搬运",深度定制转换逻辑优先写自己的函数;重负载搬运注意与 Broker 分节点部署,别让 IO 作业抢了消息收发的资源。

二、事务:跨主题原子写入的护法

订单场景有个经典两难:扣库存的结果要写库存主题,同时原始事件要标记已处理——两步跨主题,中间崩了就不一致。Pulsar 的事务把多条发送与消费确认打包成原子单元:

Transaction txn = client.newTransaction() .withTransactionTimeout(30, TimeUnit.SECONDS) .build() .get(); // 两条发送与一次消费确认,全部挂进同一事务 producerStock.newMessage().value(deductResult).transaction(txn).sendAsync(); producerAudit.newMessage().value(auditMark).transaction(txn).sendAsync(); consumer.acknowledgeAsync(msgId, txn); txn.commit().get(); // 全部生效;任一环节失败则整体回滚

事务的语义边界要背下来:它原子化的是"Pulsar 世界内"的操作——消息发送与消费确认;世界外的动作(数据库写入、外部调用)不在护法范围内,跨界一致性仍要业务层自己设计。开启事务还有存储与管理面的额外开销(事务日志、挂起状态清理),因此它留给真正需要跨主题原子的少数链路,不做全集群默认。

三、演进路标:这条水道的下一程

收拢零散的演进线索给一张路标图。部署形态:官方 Operator 与 Helm 让 Kubernetes 成为推荐部署方式,扩缩容、升级、自愈并入集群原语。架构方向:存算分离的下一层是存储进一步向对象存储扎根——分层存储正在从"冷数据卸载"演进为"数据原生住对象存储",Bookie 的本地盘角色被重新定义。查询能力:Pulsar SQL 让你能用交互式查询直接扫主题里的数据,补数与排查不必再写消费者。协议生态:Kafka 兼容协议让存量 Kafka 客户端不改代码接入,迁移期缓冲带价值巨大。

评价这些演进对自家业务的含金量,判断框架仍是第 1 章的老三样:团队形态、流量形状、未来变化——Functions 值不值、事务开不开、何时上 K8s,全都能从这三问里推出方向。

本节要点回顾

  • Functions 定位是"搬进水道的胶水代码",轻状态转换用它,复杂管道交给专业引擎;
  • 连接器分 Source 与 Sink 两个方向,标准搬运用现成的,定制逻辑写函数;
  • 事务原子化 Pulsar 世界内的发送与确认,跨界一致性仍归业务层;
  • K8s 化、存储向对象存储扎根、SQL 查询、Kafka 兼容,是明确的演进路标;
  • 一切功能取舍回到老三问:团队形态、流量形状、未来变化。

全书到此收拢:从水道演化到生产护航,一条消息的漂流全程已经走完。愿你带着这套问题框架,去读任何一张消息系统的架构图。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U