ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Pulsar Topic Compaction 实战指南:原理、自动/手动触发与消费者配置

Apache Pulsar Topic Compaction 实战指南:原理、自动/手动触发与消费者配置 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读Topic compaction主题压缩是 Apache Pulsar 提供的一项数据整理能力它允许你基于消息 key对持久化主题进行压缩让每个 key 只保留最新的一条消息从而加快对主题历史数据的读取速度。本文以 Pulsar 2.2.0 官方 cookbook 为骨架结合当前仓库gh_mirrors/pulsar28/pulsar的源码实现系统讲解 compaction 的适用场景、namespace 级自动触发策略、两种手动触发方式pulsar-admin topics compact与pulsar compact-topic、以及 Java 客户端消费 compacted topic 的完整配置。读完本文你将能独立判断业务是否适合使用 compacted topic并正确配置、触发与消费压缩后的主题。Topic compaction 是什么按 key 保留最新值Pulsar 的 topic compaction 概念说明 定义了一个核心能力通过 compaction你可以创建压缩后的主题compacted topic其中较旧的、被遮蔽obscured的条目会被修剪掉从而让读者更快地遍历主题历史。哪些消息会被视为过期/无关完全取决于你的业务场景。compaction 在 Pulsar 中是按 key 进行的消息被压缩的依据是它的 key因此使用 compaction 前必须给消息设置 key。例如股票行情场景中股票代码AAPL、GOOG、TWTR等就可以作为 key没有 key 的消息会被 compaction 过程原样保留、不参与压缩。注意compaction 只作用于带 key 的消息。可以把 key 理解为压缩所沿着的轴不带 key 的消息会被 compaction 直接忽略。底层实现两阶段压缩器从源码结构看当前仓库的 compaction 核心实现在 pulsar-broker/src/main/java/org/apache/pulsar/compaction 包下核心类是 TwoPhaseCompactor.java其类注释明确写道compaction 会分两遍two passes扫描整个主题第一遍phase one选出每个 key 在主题中的最新 offsetlatestForKey映射key → 最新 MessageId。只保存消息 ID 而不保存消息体避免把大量消息负载常驻内存第二遍phase two重新读取主题仅把该 key 的最新消息写入一个新建的 BookKeeper ledger 中形成 compacted topic 的底层存储。从 Compactor.java 可以看到压缩过程使用名为__compaction的专属订阅COMPACTION_SUBSCRIPTION压缩产物 ledger 会在元数据中打上CompactedTopicLedger属性标记。第二遍写入时还通过信号量限制在途写入数量MAX_OUTSTANDING 500控制内存占用。也正因为第一遍要把每个 key 的最新 ID 保存在内存映射里当主题的 keyspace 很大key 数量非常多时会产生内存压力——这正是下文何时用独立进程运行 compaction的依据。什么时候应该使用 compacted topic最经典的例子就是股票行情主题。假设一个主题持续接收带 key 的股票价格消息key 就是股票代码那么 compaction 之后该主题的消费者就有两种选择消费方式读取内容适用场景非压缩原始主题全部历史消息需要对过去一小时所有价格做数值计算的批量分析压缩后的主题每个 key 的最新消息实时行情大屏、最新价格展示在stock-values主题上你可以让部分消费者读取全部消息做统计分析同时让驱动实时行情的消费者读取 compacted topic避免被迫处理大量过期消息。选择哪一种由消费者的 配置 决定。一个关键优势是compaction 不会破坏原始主题。压缩过程会把原始主题原封不动地保留相当于在原始主题旁新增了一个压缩变体。因此你可以放心地对某个主题执行 compaction需要访问非压缩版本的消费者完全不受影响。配置自动压缩namespace 级策略租户管理员tenant administrator可以在namespace 级别配置 compaction 策略策略指定主题积压backlog增长到多大时自动触发压缩。该策略对该 namespace 下的所有主题生效。例如当积压达到 100MB 时触发压缩$ bin/pulsar-admin namespaces set-compaction-threshold \ --threshold 100M my-tenant/my-namespace在 2.2.0 这个版本里compaction threshold 是按 namespace 配置并作用于其下所有主题的。从当前仓库的管理接口实现看这一配置链路已经演进得更加完整CmdTopics.java 中提供了按主题级别操作的三个命令# 查看主题当前的压缩阈值--applied 表示是否展示最终生效值 $ bin/pulsar-admin topics get-compaction-threshold --applied \ persistent://my-tenant/my-namespace/my-topic # 设置主题压缩阈值格式如 10M / 16G / 3T0 表示禁用自动压缩 $ bin/pulsar-admin topics set-compaction-threshold \ --threshold 10M persistent://my-tenant/my-namespace/my-topic # 移除主题的压缩阈值 $ bin/pulsar-admin topics remove-compaction-threshold \ persistent://my-tenant/my-namespace/my-topic自动触发的实现细节与 broker 配置项从 PersistentTopic.java 的源码可以看到自动压缩的判定逻辑broker 周期性检查每个带压缩策略的主题计算__compaction订阅对应的积压估算值estimateBacklogSize一旦积压超过阈值就调用triggerCompaction()。同时如果主题配置了压缩策略broker 会确保__compaction游标存在为后续压缩做好准备。与自动压缩相关的 broker 配置项定义在 ServiceConfiguration.java对应的 conf/broker.conf 示例如下# 检查带压缩策略的主题是否需要压缩的间隔秒默认 60 brokerServiceCompactionMonitorIntervalInSeconds60 # 全局默认压缩阈值字节0 表示默认禁用自动压缩默认 0 brokerServiceCompactionThresholdInBytes0 # 压缩第一阶段循环的超时秒若第一阶段执行时间超过该值压缩将不会继续默认 30 brokerServiceCompactionPhaseOneLoopTimeInSeconds30说明namespace/主题级阈值生效的优先级高于该全局默认值从 PersistentTopicsBase.java 的实现看生效阈值依次取主题策略 → namespace 策略 → broker 全局配置。手动触发压缩的两种方式方式一通过管理 API推荐使用pulsar-adminCLI 工具的topics compact命令例如$ bin/pulsar-admin topics compact \ persistent://my-tenant/my-namespace/my-topic该命令通过 Pulsar 的 REST 管理 API 触发压缩。从 CmdTopics.java 可以看到命令最终调用getTopics().triggerCompaction(persistentTopic)在 broker 侧PersistentTopicsBase.java 中的internalTriggerCompaction会校验 namespace 归属后调用((PersistentTopic) topic).triggerCompaction()。另外还提供了查询压缩运行状态的命令# 查询主题压缩状态--wait 表示等待压缩完成 $ bin/pulsar-admin topics compaction-status \ --wait persistent://my-tenant/my-namespace/my-topic该命令会输出Compaction has not been run未执行过、Compaction is currently running运行中、Compaction was a success成功或报错信息对应实现见 CmdTopics.java。方式二独立进程运行pulsar compact-topic如果不希望通过 REST API 在 broker 进程内运行压缩可以使用pulsar compact-topic命令$ bin/pulsar compact-topic \ --topic persistent://my-tenant/my-namespace/my-topicpulsar compact-topic会直接与 ZooKeeper 通信因此pulsarCLI 工具需要一份有效的 broker 配置。配置默认放在conf/broker.conf如果配置在别处可以用--broker-conf指定# 指定非默认位置的 broker 配置 $ bin/pulsar compact-topic \ --broker-conf /path/to/broker.conf \ --topic persistent://my-tenant/my-namespace/my-topic # 配置就在默认的 conf/broker.conf 时无需额外参数 $ bin/pulsar compact-topic \ --topic persistent://my-tenant/my-namespace/my-topic该命令的实现入口是 CompactorTool.java-c/--broker-conf指定配置文件-t/--topic指定要压缩的主题最终调用compactor.compact(topic)并输出压缩产物 ledger 的 IDCompaction of topic ... complete. Compacted to ledger ...。什么时候推荐用独立进程运行压缩当你希望避免干扰 broker 性能时。不过只有对 keyspace 很大的主题key 数量非常多做压缩时broker 性能才可能受影响——因为压缩第一遍会在内存中为每个 key 保留一份副本key 数量越多内存压力越大。在绝大多数情况下通过pulsar-admin topics compact走 REST API 在 broker 内运行压缩不会有什么问题使用pulsar compact-topic应被视为一种边缘场景。压缩频率如何把握多久触发一次压缩没有统一答案完全取决于业务场景。如果你希望 compacted topic 的读取极其快速那就应该相当频繁地运行压缩让每个 key 保留的最新值尽量靠近当前时刻。消费者配置如何读取 compacted topicPulsar 的消费者consumer和读取者reader需要显式配置才能读取压缩后的主题。各语言客户端的开启方式如下。Java 客户端Java consumer 读取 compacted topic 需要把readCompacted设置为trueConsumerbyte[] compactedTopicConsumer client.newConsumer() .topic(some-compacted-topic) .readCompacted(true) .subscribe();从当前仓库的客户端 API 定义看ConsumerBuilder.java 中readCompacted(boolean)的 Javadoc 明确指出readCompacted只能用于持久化主题persistent topic的订阅且要求该订阅只有单一活跃消费者同一订阅下有多个消费者竞争时无法保证读取压缩视图。与之对应ReaderBuilder.java 也为 reader 提供了同样的readCompacted开关。由于 compaction 按 key 进行你在 compacted topic 上生产的消息必须携带 keykey 的内容取决于业务如股票代码、用户 ID、设备 ID 等。没有 key 的消息会被压缩过程忽略。下面是带 key 消息的构造示例import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageBuilder; Messagebyte[] msg MessageBuilder.create() .setContent(someByteArray) .setKey(some-key) .build();完整的生产带 key 消息到 compacted topic示例import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageBuilder; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; PulsarClient client PulsarClient.builder() .serviceUrl(pulsar://localhost:6650) .build(); Producerbyte[] compactedTopicProducer client.newProducer() .topic(some-compacted-topic) .create(); Messagebyte[] msg MessageBuilder.create() .setContent(someByteArray) .setKey(some-key) .build(); compactedTopicProducer.send(msg);如果消费者没有开启 readCompacted 会怎样只要没有开启readCompacted消费者就仍然可以读取未压缩的原始主题行为与普通主题完全一致。也就是说readCompacted只是切换消费者看到的视图false默认读取全量历史true读取压缩后每个 key 的最新值。小结与使用检查清单环节关键操作/配置依据生产端消息必须设置 keysetKey无 key 消息不参与压缩本文生产带 key 消息示例自动触发bin/pulsar-admin namespaces set-compaction-threshold --threshold 100M ns2.2.0 按 namespace 生效cookbook 原文手动触发bin/pulsar-admin topics compact topicREST API推荐CmdTopics.java独立进程压缩bin/pulsar compact-topic --topic topic必要时加--broker-confCompactorTool.java消费端JavareadCompacted(true)仅限持久化主题且订阅单一活跃消费者ConsumerBuilder.javabroker 调优brokerServiceCompactionMonitorIntervalInSeconds、brokerServiceCompactionThresholdInBytes、brokerServiceCompactionPhaseOneLoopTimeInSecondsconf/broker.conf使用 compaction 的核心心智模型是key 是压缩的轴compaction 保留每个 key 的最新消息原始主题始终不变。给消息打上合适的 key选择自动或手动触发再让需要的消费者打开readCompacted即可在同一份数据上同时支撑全量历史分析与最新状态读取两种截然不同的消费模式。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Topic Compaction 实战指南原理、触发方式与消费者配置Apache Pulsar Topic Compaction 实战指南原理、触发方式与消费者配置 本文围绕 Pulsar 官方 Cookbook 文档 coo消息队列后端流处理Apache Pulsar Topic Compaction 实战指南配置、触发与消费者读取Apache Pulsar Topic Compaction 实战指南配置、触发与消费者读取 本指南基于 Pulsar 官方文档 cookbooks comp消息队列后端流处理Apache Pulsar Topic Compaction 完全指南原理、触发方式与消费者配置Apache Pulsar Topic Compaction 完全指南原理、触发方式与消费者配置 Apache Pulsar 的 Topic Compacti消息队列后端流处理上一篇Unitree Go2 ROS2 SDK深度技术解析从架构设计到实战应用下一篇暗黑破坏神2存档修改器终极指南从新手到高手的完整解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表