本节摘要:远洋水域有三样装备:Functions 把消费逻辑简化成一段部署进集群的代码,连接器把外部系统接进水道,事务让跨主题操作原子化;再把视线放远,Kubernetes 化部署与存算进一步解耦是这条水道明确的下一程。本节是生态全景加演进路标,帮你在功能清单与团队承接力之间做出清醒取舍。
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)分两个方向: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,全都能从这三问里推出方向。
全书到此收拢:从水道演化到生产护航,一条消息的漂流全程已经走完。愿你带着这套问题框架,去读任何一张消息系统的架构图。