ARTICLE DETAIL

资讯详情

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

Apache Pulsar 主题压缩(Topic Compaction)实战指南:原理、配置与客户端接入

Apache Pulsar 主题压缩(Topic Compaction)实战指南:原理、配置与客户端接入 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本文基于 Apache Pulsar 2.3.0 版本文档与仓库源码系统讲解Topic Compaction主题压缩的核心机制与实操方法从按 key 保留最新消息的工作原理出发覆盖命名空间级自动压缩策略、pulsar-admin/pulsar compact-topic两种手动触发方式、Broker 相关配置参数以及 Java 客户端的生产与消费接入。读完本文你将能够独立为一个以 key 为维度、只需最新状态的业务场景如股票行情、设备状态快照设计并落地压缩主题。什么是 Topic CompactionPulsar 的 Topic Compaction 特性允许你创建compacted topics压缩主题在压缩过程中主题中较旧的、已被遮蔽obscured的条目会被裁剪从而加快读者回溯主题历史时的读取速度哪些消息被视为过时/无关取决于你的具体业务场景。压缩以per-key basis按 key 维度进行——消息根据其 key 被压缩。例如在股票行情场景中AAPL、GOOG这类股票代码可作为消息的 key没有 key 的消息会被压缩过程原样保留、不做处理。key 因此可以理解为压缩作用的轴axis。从源码结构看压缩的核心实现位于 Compactor.java其中定义了专用的内部订阅COMPACTION_SUBSCRIPTION __compactionCompactor.java#L36以及标识压缩后 ledger 的元数据属性CompactedTopicLedger压缩流程通过RawReader读取原始消息并最终写回一个新的 ledger实现原主题不动、另建压缩视图的效果。什么时候应该使用压缩主题压缩主题的典型场景是股票行情类主题生产者以股票代码GOOG、AAPL、TWTR等作为消息 key 持续推送价格数据。压缩之后订阅同一主题的消费者拥有两种选择读取原始未压缩主题获取历史数据即主题中的全部消息例如用于对最近一小时的所有价格做批量计算读取压缩主题只想看到每个 key 的最新消息例如驱动实时行情展示避免被迫处理大量过时消息。以名为stock-values的主题为例部分消费者可以订阅未压缩版本处理全量数据而支撑实时行情的消费者则通过设置readCompacted读取压缩版本。消费者究竟拉取哪个变体完全由消费者的客户端配置决定。这里有一个关键优势压缩过程不会破坏原始主题它本质上新增了一条可选的替代读取路径。你可以对一个主题执行压缩而需要未压缩版本的消费者完全不受影响无需在两者之间做取舍。配置自动压缩命名空间级策略租户管理员可以在namespace命名空间级别配置压缩策略指定主题 backlog 增长到多大时触发压缩。例如设置 backlog 达到 100MB 时触发压缩$ bin/pulsar-admin namespaces set-compaction-threshold \ --threshold 100M my-tenant/my-namespace命名空间上的压缩阈值配置会应用到该命名空间下的所有主题。对应的管理 REST API 位于 Namespaces.javaGET /{tenant}/{namespace}/compactionThreshold查询阈值PUT /{tenant}/{namespace}/compactionThreshold设置阈值描述为在压缩被触发前主题中允许的最大未压缩字节数。Broker 侧相关配置除了命名空间策略Broker 还提供三个与压缩调度直接相关的配置项见 conf/broker.conf配置项默认值说明brokerServiceCompactionMonitorIntervalInSeconds60Broker 周期性检查带有压缩策略的主题是否需要压缩的间隔秒。brokerServiceCompactionThresholdInBytes0当估算的 backlog 大小超过该阈值时触发压缩设为0表示禁用该检查。brokerServiceCompactionPhaseOneLoopTimeInSeconds30压缩第一阶段单轮循环的超时时间超过该时间则压缩不再继续。其中brokerServiceCompactionPhaseOneLoopTimeInSeconds在 TwoPhaseCompactor.java 中被转换为phaseOneLoopReadTimeout使用用于约束两阶段压缩中扫描并记录每个 key 最新消息这一阶段的单轮执行时长。手动触发压缩方式一通过 pulsar-admin 管理 API运行压缩最直接的方式是使用pulsar-adminCLI 工具的topics compact命令$ bin/pulsar-admin topics compact \ persistent://my-tenant/my-namespace/my-topicpulsar-admin通过 Pulsar 的 REST API 触发压缩。对应实现位于 PersistentTopicsBase.javaBroker 会基于命名空间上配置的压缩阈值自动判断是否需要执行压缩。方式二使用独立的 compact-topic 进程如果希望压缩在独立进程中运行而非经过 REST API可以使用pulsar compact-topic命令$ bin/pulsar compact-topic \ --topic persistent://my-tenant/my-namespace/my-topic该命令直接与 ZooKeeper 通信因此需要读取 Broker 配置来获知集群元数据地址。默认使用conf/broker.conf若配置位于其他路径可通过--broker-conf显式指定$ bin/pulsar compact-topic \ --broker-conf /path/to/broker.conf \ --topic persistent://my-tenant/my-namespace/my-topic从源码看CompactorTool.java 定义了该工具的完整参数-c/--broker-conf默认指向当前目录下conf/broker.conf、-t/--topic必填待压缩主题、-h/--help与-g/--generate-docs。工具的进程内压缩逻辑在 TwoPhaseCompactor.java 中实现分为两个阶段第一阶段扫描主题、在内存中为每个 key 记录最新消息的 IDMapString, MessageId latestForKey第二阶段基于该映射把每个 key 的最新消息写入新的压缩 ledger。何时使用独立进程当主题的 keyspace 很大即 key 数量非常多时压缩第一阶段会在内存中为每个 key 保留一条记录可能造成内存压力此时建议用独立进程避免影响 Broker 性能。绝大多数场景下pulsar-admin topics compactREST API 方式不会出现问题pulsar compact-topic应被视为边界场景。触发频率建议压缩触发的频率应结合具体场景权衡。如果你希望压缩主题的读取速度极快则应相对频繁地运行压缩反之读取频率较低的场景可以放宽压缩间隔。消费者配置如何读取压缩主题Pulsar 的消费者Consumer与读取器Reader需要显式配置才能从压缩主题读取。如果未做该配置消费者仍可正常读取未压缩的主题只是拿不到仅每个 key 最新消息的视图。需要注意readCompacted只能用于**持久主题上、单活跃消费者即 Failover 或 Exclusive 订阅模式**的订阅对非持久主题或 Shared 订阅启用该选项会导致订阅调用抛出PulsarClientException见 ConsumerBuilder.java。Java 客户端使用 Java 消费者读取压缩主题时需将readCompacted参数设为trueConsumerbyte[] compactedTopicConsumer client.newConsumer() .topic(some-compacted-topic) .readCompacted(true) .subscribe();由于压缩按 key 进行生产到压缩主题上的消息必须携带 keykey 的具体内容取决于业务如股票代码。没有 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();完整的生产示例连接本地单机 Pulsar 并向压缩主题发送带 key 消息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);总结Topic Compaction 为 Pulsar 提供了状态快照式的读取视角生产者只需为消息设置 keyBroker 通过命名空间阈值或手动命令在后台将每个 key 收敛为最新一条消息原始主题则保持完整不变。落地时只需三步为消息加 key→配置/触发压缩→在消费者上开启readCompacted。无论是股票行情、设备状态、配置快照等读最新类场景压缩主题都能在不影响历史数据消费者的情况下显著提升最新值读取的效率。延伸阅读压缩概念总览concepts-topic-compaction.mdpulsar-admin完整命令参考reference-pulsar-admin.mdCLI 工具参考含pulsar compact-topicreference-cli-tools.mdBroker 配置参考reference-configuration.md压缩核心实现Compactor.java、TwoPhaseCompactor.java赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Topic Compaction 详解原理、配置与实战Apache Pulsar Topic Compaction 详解原理、配置与实战 Topic Compaction 是 Apache Pulsar 提供的一消息队列后端流处理Apache Pulsar Topic Compaction 全面解析原理、配置与实战Apache Pulsar Topic Compaction 全面解析原理、配置与实战 Pulsar 的 Topic Compaction主题压缩 是构建消息队列后端流处理Apache Pulsar 主题压缩Topic Compaction完整实战指南从自动/手动触发到消费者接入Apache Pulsar 主题压缩Topic Compaction完整实战指南从自动/手动触发到消费者接入 主题压缩Topic Compaction消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表