Pulsar调优与实践
Apache Pulsar 性能优化指南
Apache Pulsar 是一个分布式消息系统,设计用于处理大规模、高吞吐量、低延迟的实时数据流。要对其性能进行优化,可以从以下几个关键方面入手:
1. 硬件配置与集群规模
硬件资源
- CPU、内存和磁盘I/O:确保 Broker 和 BookKeeper 节点具有足够的资源。根据预计的吞吐量和消息大小,选择合适的硬件配置,特别是使用高速 SSD 以减少磁盘访问延迟。
- 网络带宽:保证集群内部以及客户端与集群之间的网络带宽充足且稳定,减少网络瓶颈影响。
- 集群规模:根据实际负载动态调整 Broker 和 BookKeeper 节点的数量,以分散负载并提高容错能力。合理设置分区数和副本数,以平衡可用性、持久性和性能。
2. 配置参数调优
BookKeeper
- writeQuorumSize、ackQuorumSize、journalMaxSizeMB:调整这些配置以适应特定的工作负载和可靠性要求。
Broker
- maxConcurrentLookupRequest、maxConcurrentNamespaceBundleSplit、systemTopicEnabled:调整这些配置以优化服务端处理能力和资源利用率。
客户端
- batchingEnabled、batchingMaxMessages、batchingMaxPublishDelayMs:调整这些参数以减少网络交互次数和提升发送/接收效率。
3. 负载均衡
Namespace 级别
- 利用 Pulsar 的多租户特性,合理划分 Namespace,分配到不同的 Broker 上,实现负载均衡。
Topic 级别
- 通过 Pulsar 的自动或手动 Topic 分区重分布功能,调整 Topic 分区在 Broker 间的分布,避免热点。
Consumer 组
- 优化订阅模式(Exclusive、Shared、Failover、Key_Shared 等)和负载均衡策略(RoundRobin、LeastBacklog、MessageRateSubscribe 等),以更高效地分发消息。
4. 数据存储与压缩
存储格式
- 评估不同存储格式(如 dbLedgerStorage)的性能特点,根据实际需求进行选择或优化。
消息压缩
- 启用消息压缩(如 LZ4、Snappy、ZLIB 等),减少网络传输和存储空间占用,特别是在带宽有限或存储成本敏感的情况下。
5. 监控与分析
监控系统
- 部署完善的监控系统,跟踪 Broker、BookKeeper、客户端的各项性能指标(如 CPU 使用率、内存使用、磁盘 I/O、网络流量、队列深度、延迟等),及时发现性能瓶颈。
日志分析
- 定期审查系统日志,识别异常行为和潜在问题。使用工具如 Pulsar Manager 或 Prometheus/Grafana 等进行可视化监控和报警设置。
性能测试
- 使用工具如
pulsar-perf进行基准测试和压力测试,量化性能改进效果,为调优决策提供依据。
6. 高级特性利用
TieredStorage
- 启用分层存储,将长期保存的数据自动移至低成本存储(如 S3、HDFS 等),释放 BookKeeper 资源,降低存储成本并提高整体性能。
SchemaRegistry
- 使用 Schema Registry 管理消息结构,减少数据冗余,提高序列化/反序列化效率。
FunctionWorker
- 如果使用 Pulsar Functions,优化 Function Worker 配置,如并行度、资源分配等,以提升计算性能。
7. 运维最佳实践
定期维护
- 执行定期的 BookKeeper ledger 清理、Broker 垃圾回收、Namespace Bundle 分裂等操作,保持系统健康状态。
故障恢复
- 制定详尽的故障恢复流程和预案,确保在节点故障时能够快速恢复服务。
版本升级
- 关注 Pulsar 社区的新版本发布,适时升级以获取性能优化和 bug 修复。
8. 消息 TTL、Retention、消息配额
基于业务周期
- 根据业务处理周期确定消息的最大存活时间。例如,如果业务逻辑要求消息在接收到后的 2 小时内处理完毕,那么 TTL 可设置为略大于 2 小时。
存储成本与性能权衡
- 在满足业务需求的前提下,尽可能缩短 TTL 以减少存储成本和清理过期消息带来的系统开销。但也要避免设置过短的 TTL 导致合法消息被提前清除。
监控与调整
- 实施 TTL 设置后,密切关注系统监控数据,如存储使用情况、消息清理速率等。根据实际情况动态调整 TTL 值,确保既能满足业务需求,又能有效控制存储成本和系统资源。
设置保留策略
- 命名空间内的每个主题的大小限制设置为 10 GB,时间限制设置为 3 小时(即三小时内最大限制为 10GB)。
bash
pulsar-admin namespaces set-retention my-tenant/my-ns
--size 10G
--time 3h
- 时间不受限制,大小限制设置为 1 TB。大小限制决定了保留。
bash
pulsar-admin namespaces set-retention my-tenant/my-ns
--size 1T
--time -1
- 要实现无限保留,请将两个值都设置为 -1。
bash
pulsar-admin namespaces set-retention my-tenant/my-ns
--size -1
--time -1
- 要禁用保留策略,请将两个值都设置为 0。
bash
pulsar-admin namespaces set-retention my-tenant/my-ns
--size 0
--time 0
清除积压
bash
pulsar-admin namespaces clear-backlog my-tenant/my-ns
TTL
- 为命名空间设置 TTL
bash
pulsar-admin namespaces set-message-ttl my-tenant/my-ns
--messageTTL 120
- 获取命名空间的 TTL 配置
bash
pulsar-admin namespaces get-message-ttl my-tenant/my-ns
- 删除命名空间的 TTL 配置
bash
pulsar-admin namespaces remove-message-ttl my-tenant/my-ns
设置 Backlog Quota
bash
pulsar-admin topics set-backlog-quota
--topic persistent://tenant/namespace/topic
--limit <limit-type>
[--size <size-in-bytes>|--count <count-of-messages>]
[--policy <policy-name>]
参数说明
--topic: 指定要设置 backlog quota 的 Topic 全路径,格式为persistent://tenant/namespace/topic。--limit: 指定限制类型,可选值为destination_storage或message_age。destination_storage: 基于积压消息的总存储大小进行限制。message_age: 基于积压消息的最老消息的存活时间(age)进行限制。注意,此选项需要 broker 版本 >= 2.8.0。
--size: 当--limit destination_storage时,指定积压消息的最大存储大小(单位:字节)。例如:--size 1073741824(即 1 GB)。--count: 当--limit message_age时,指定积压消息的最大数量。例如:--count 1000000。--policy: (可选)指定 quota 超限后的处理策略。可选值为:producer_request_hold: 当积压消息达到配额时,暂停生产者发送请求,直到消费者开始消费并减少积压。producer_exception: 当积压消息达到配额时,立即抛出异常给生产者,阻止其继续发送消息。consumer_backlog_eviction: 当积压消息达到配额时,自动删除最旧的消息,以腾出空间给新消息。如果不指定,默认策略为producer_request_hold。
基于存储大小限制
bash
pulsar-admin topics set-backlog-quota
--topic persistent://my-tenant/my-namespace/my-topic
--limit destination_storage
--size 1073741824
基于消息数量限制
bash
pulsar-admin topics set-backlog-quota
--topic persistent://my-tenant/my-namespace/my-topic
--limit message_age
--count 1000000
--policy consumer_backlog_eviction
9. 配置参数调优
Broker 配置
核心参数
brokerServicePort / brokerServicePortTls: 设置 Broker 的监听端口,确保客户端能够正确连接。tickDurationSeconds: Broker 心跳间隔,影响 Broker 内部定时任务的执行频率。webServicePort / webServicePortTls: 设置 Broker 的管理接口端口,用于对接 Prometheus、Grafana 等监控工具。globalZookeeperServers: 设置全局 ZooKeeper 集群地址,用于服务发现和协调。configurationStoreServers: 设置配置存储 ZooKeeper 集群地址,用于存储和管理 Pulsar 的元数据。loadBalancerEnabled: 控制是否启用自动负载均衡。根据业务需求,可以设置为false禁止自动均衡,或者调整相关负载均衡策略。systemTopicEnabled: 如果系统中使用了 Pulsar Functions 或 Schema Registry 等依赖系统 Topic 的功能,应保持为true。否则,如果不需要这些功能,可设为false以减少资源消耗。maxConcurrentLookupRequest: 设置 Broker 同时处理的查找请求的最大数量。根据客户端数量和消息查找频率适当调整。maxConcurrentNamespaceBundleSplit: 控制并发进行的 Namespace Bundle 分裂任务数量。根据集群规模和 Namespace 管理需求设定。
性能与资源管理
systemResourceUsageCheckIntervalSeconds: 控制 Broker 检查系统资源使用情况的间隔。maxConcurrentLookupRequests: 设置 Broker 同时处理的查找请求的最大数量,根据客户端数量和消息查找频率适当调整。maxConcurrentNonPersistentMessagePerConnection: 控制每个连接上并发处理的非持久化消息数量。maxPendingPublishRequestsPerConnection: 控制每个连接上待处理的发布请求最大数量。maxUnackedMessagesPerConsumer: 控制每个消费者未确认消息的最大数量,防止消费者过载。maxUnackedMessagesPerSubscription: 控制每个订阅未确认消息的最大数量,防止订阅过载。messageChunkingEnabled: 是否开启消息分块。对于大消息场景,开启分块有助于减少内存压力和网络传输负担。chunkingMaxMessageSize: 分块消息的最大大小。应与客户端配置保持一致。maxMessageSize: Broker 允许接收的最大单个消息大小。根据业务消息大小设置。nettyMaxFrameSizeBytes: Netty 的最大帧大小。maxConcurrentConnections: Broker 允许的最大并发连接数。ServiceNumWorkerThreads: 该参数指定了每个 broker 节点处理请求的工作线程数量。增加该值可以提高 broker 的并发处理能力,但同时也会增加系统资源的消耗。默认值为 16。ServiceNumIoThreads: 该参数指定了每个 broker 节点处理网络 I/O 的线程数量。增加该值可以提高 broker 的网络吞吐量,但同时也会增加系统资源的消耗。默认值为 8。ServiceNumSchedulerThreads: 该参数指定了每个 broker 节点处理定时任务的线程数量。增加该值可以提高 broker 的定时任务处理能力,但同时也会增加系统资源的消耗。默认值为 8。ServiceNumExecutorThreadPoolSize: 该参数指定了每个 broker 节点处理请求的线程池大小。增加该值可以提高 broker 的并发处理能力,但同时也会增加系统资源的消耗。默认值为 20。ServiceNumListenerThreads: 该参数指定了每个 broker 节点处理网络连接的线程数量。增加该值可以提高 broker 的网络连接处理能力,但同时也会增加系统资源的消耗。默认值为 1。
BookKeeper 交互
managedLedgerDefaultEnsembleSize, managedLedgerDefaultWriteQuorum, managedLedgerDefaultAckQuorum: 控制 BookKeeper ledger 的写入和确认策略,根据可用性和性能需求调整。managedLedgerMaxEntriesPerLedger: 每个 ledger 的最大条目数,影响 BookKeeper 的滚动频率和数据分布。managedLedgerCursorRolloverTimeInSeconds: 游标自动滚动的时间阈值,影响过期消息清理速度。managedLedgerCacheSizeMB: 设置 BookKeeper 缓存大小,影响读取性能。managedLedgerMinLedgerRolloverTimeMinutes: 控制 ledger 最小滚动间隔,防止频繁滚动。
其他参数
authenticationProviders: 配置支持的认证提供者,如 Athenz、JWT、TLS 等。authorizationProviders: 配置支持的授权提供者,如 Pulsar Authorization Provider、Athenz 等。tlsCertificateFilePath, tlsKeyFilePath: 设置 Broker 的 TLS 证书和私钥路径,启用 TLS 加密。superUserRoles: 设置超级用户角色,拥有所有权限。
BookKeeper 配置
核心参数
zkServers: 设置 ZooKeeper 集群地址,确保 BookKeeper 能正确连接。ledgersPath: BookKeeper 元数据在 ZooKeeper 中的存储路径。journalDirectories: 设置 Bookie 的 Journal 存储目录。ledgerDirectories: 设置 Bookie 的 Ledger 存储目录。
性能与资源管理
bookieReadBufferSizeBytes: 设置 Bookie 读取缓冲区大小,影响读取性能。numAddWorkerThreads: 设置 Bookie 用于处理写入请求的线程数。numReadWorkerThreads: 设置 Bookie 用于处理读取请求的线程数。journalMaxSizeMB: 控制每个 Bookie 的 Journal 文件最大大小,影响数据刷盘策略。journalBufferedWritesThreshold: 控制 Journal 缓冲区写入阈值,影响写入性能。journalFlushWhenQueueEmpty: 控制是否在写入队列为空时强制刷盘,对数据一致性有重要影响。journalAdaptiveGroupWrites: 是否启用 Journal 自适应组写,优化写入性能。diskUsageThreshold: 磁盘使用率达到此阈值时触发警告。gcWaitTime: 控制 BookKeeper Garbage Collection 等待时间。
其他参数
tlsEnable: 是否启用 TLS 加密。tlsClientAuthentication: 是否要求客户端进行 TLS 认证。tlsKeyStorePath, tlsKeyStorePasswordPath: 设置 BookKeeper 的 TLS 证书库路径和密码文件路径。tlsTrustStorePath, tlsTrustStorePasswordPath: 设置 BookKeeper 的 TLS 信任证书库路径和密码文件路径。
客户端配置
生产者
batchingEnabled: 是否开启消息批处理,通常开启可以提高发送效率。batchingMaxMessages: 批处理中包含的最大消息数。batchingMaxPublishDelayMs: 批处理最大等待时间,超过此时间即使未达到batchingMaxMessages也会发送。batchingMaxAllowedSizeInBytes: 批处理的最大字节数。batchingPartitionSwitchFrequencyByPublishDelay: 控制因发布延迟而切换分区的频率。sendTimeoutMs: 生产者发送消息的超时时间。blockIfQueueFull: 当发送队列满时是否阻塞生产者。compressionType: 选择压缩算法(如 LZ4、ZLIB、SNAPPY 等)。compressionLevel: 设置压缩等级(如适用的话)。compressionThreshold: 设置触发压缩的最小消息大小。
消费者
receiverQueueSize: 消费者接收队列大小,影响消费速度和内存使用。subscriptionType: 订阅类型(Exclusive、Shared、Failover、Key_Shared 等),根据消费模式选择。negativeAckRedeliveryDelayMs: 负面确认后重新投递消息的延迟。ackTimeoutMillis: 自动 ACK 超时时间,影响消息确认策略。numIOThreads: 设置 Netty 工作线程数,一般根据 CPU 核心数进行调整。numWorkerThreads: 设置 Pulsar 内部工作线程数,用于处理非 I/O 密集型任务。MaxTotalReceiverQueueSizeAcrossPartitions: 该参数指定了消费者在所有分区上允许的最大接收队列大小。增加该值可以提高消费者的并发性能,但同时也会增加内存的消耗。默认值为 50000。MaxPendingChuckedMessage: 该参数指定了消费者允许等待确认的分块消息数量。增加该值可以提高消费者的并发性能,但同时也会增加内存的消耗。默认值为 1000。
10. Flink 与 Pulsar 集成以及相关优化
Pulsar Source
使用示例
java
PulsarSource<String> source = PulsarSource.builder()
.setConfig(PulsarSourceOptions.PULSAR_PARTITION_DISCOVERY_INTERVAL_MS, 10000)
.setServiceUrl(serviceUrl)
.setAdminUrl(adminUrl)
.setStartCursor(StartCursor.earliest())
.setTopics("my-topic")
.setDeserializationSchema(new SimpleStringSchema())
.setSubscriptionName("my-subscription")
.build();
env.fromSource(source, WatermarkStrategy.noWatermarks(), "Pulsar Source");
反序列化器
使用 Pulsar 的 Schema 解析消息。如果使用 KeyValue 或者 Struct 类型的 Schema,那么 Pulsar 的 Schema 将不会含有类型类信息,但 PulsarSchemaTypeInformation 需要通过传入类型类信息来构造。因此我们提供的 API 支持用户传入类型信息。
java
// 基础数据类型
PulsarSourceBuilder.setDeserializationSchema(Schema);
// 结构类型 (JSON, Protobuf, Avro, etc.)
PulsarSourceBuilder.setDeserializationSchema(Schema, Class);
// 键值对类型
PulsarSourceBuilder.setDeserializationSchema(Schema, Class, Class);
使用 Flink 的 DeserializationSchema 解析消息。
java
PulsarSourceBuilder.setDeserializationSchema(DeserializationSchema);
使用 Flink 的 TypeInformation 解析消息。
java
PulsarSourceBuilder.setDeserializationSchema(TypeInformation, ExecutionConfig);
在 Source 中启用 Schema Evolution
java
Schema<SomePojo> schema = Schema.AVRO(SomePojo.class);
PulsarSource<SomePojo> source = PulsarSource.builder()
...
.setDeserializationSchema(schema, SomePojo.class)
.enableSchemaEvolution()
.build();
使用动态 Schema
Pulsar 提供了 Schema.AUTO_CONSUME() 这个 Schema 来消费消息。它支持所有的 Topic,并常常用于消费 Topic 中存在多种 Schema 的情况。使用此 Schema 消费时,Pulsar 会将消息解析为 GenericRecord。
但是 PulsarSourceBuilder.setDeserializationSchema(Schema) 等方法并不支持 Schema.AUTO_CONSUME() 类型。所以,我们提供了 GenericRecordDeserializer 接口来让用户定义如何解析 GenericRecord。
java
GenericRecordDeserializer<SomePojo> deserializer = ...;
PulsarSource<SomePojo> source = PulsarSource.builder()
...
.setDeserializationSchema(deserializer)
.build();
当前动态 Schema 只支持 AVRO,JSON 和 Protobuf 类型的消息解析。
消息消费游标
java
// 从 Topic 里面最早的一条消息开始消费。
StartCursor.earliest();
// 从 Topic 里面最新的一条消息开始消费。
StartCursor.latest();
// 从给定的消息开始消费。
StartCursor.fromMessageId(MessageId);
// 与前者不同的是,给定的消息可以跳过,再进行消费。
StartCursor.fromMessageId(MessageId, boolean);
// 从给定的消息发布时间开始消费,这个方法因为名称容易导致误解现在已经不建议使用。你可以使用方法 StartCursor.fromPublishTime(long)。
StartCursor.fromMessageTime(long);
// 从给定的消息发布时间开始消费。
StartCursor.fromPublishTime(long);
Pulsar Sink
使用示例
java
DataStream<String> stream = ...;
PulsarSink<String> sink = PulsarSink.builder()
.setServiceUrl(serviceUrl)
.setAdminUrl(adminUrl)
.setTopics("topic1")
.setSerializationSchema(new SimpleStringSchema())
.setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
.build();
stream.sinkTo(sink);
序列化器
序列化器(PulsarSerializationSchema)负责将 Flink 中的每条记录序列化成 byte 数组,并通过网络发送至指定的写入 Topic。和 Pulsar Source 类似的是,序列化器同时支持使用基于 Flink 的 SerializationSchema 接口实现序列化器和使用 Pulsar 原生的 Schema 类型实现的序列化器。不过序列化器并不支持 Pulsar 的 Schema.AUTO_PRODUCE_BYTES()。
如果不需要指定 Message 接口中提供的 key 或者其他的消息属性,可以从上述 2 种预定义的 PulsarSerializationSchema 实现中选择适合需求的一种使用。
使用 Pulsar 的 Schema 来序列化 Flink 中的数据。
java
// 原始数据类型
PulsarSinkBuilder.setSerializationSchema(Schema);
// 有结构数据类型(JSON、Protobuf、Avro 等)
PulsarSinkBuilder.setSerializationSchema(Schema, Class);
// 键值对类型
PulsarSinkBuilder.setSerializationSchema(Schema, Class, Class);
使用 Flink 的 SerializationSchema 来序列化数据。
java
PulsarSinkBuilder.setSerializationSchema(SerializationSchema);
在 Sink 中启用 Schema Evolution
同时使用 Pulsar 的 Schema 以及在 builder 中指定 PulsarSinkBuilder.enableSchemaEvolution() 可以启用 Schema evolution 特性。该特性会使用 Pulsar Broker 端提供的 Schema 版本兼容性检测以及 Schema 版本演进。
java
Schema<SomePojo> schema = Schema.AVRO(SomePojo.class);
PulsarSink<String> sink = PulsarSink.builder()
...
.setSerializationSchema(schema, SomePojo.class)
.enableSchemaEvolution()
.build();
如果想要使用 Pulsar 原生的 Schema 序列化消息而不需要 Schema Evolution 特性,那么写入的 Topic 会使用 Schema.BYTES 作为消息的 Schema。Pulsar 并不会存储 Schema.BYTES,所以通过此方式写入的 Topic 可能不存在 Schema 信息,对应 Topic 的消费者需要自己负责反序列化的工作。
例如,如果使用 PulsarSinkBuilder.setSerializationSchema(Schema.STRING) 而不使用 PulsarSinkBuilder.enableSchemaEvolution()。那么在写入 Topic 中所记录的消息 Schema 将会是 Schema.BYTES。
PulsarMessage <byte[]> 类型的消息的校验
Pulsar 的 topic 至少会包含一种 Schema,Schema.BYTES 是默认的 Schema 类型并常作为没有 Schema 的 topic 的 Schema 类型。使用 Schema.BYTES 发送消息将会跳过类型检测,这意味着使用 SerializationSchema 和没有启用 Schema evolution 的 Schema 所发送的消息并不安全。
可以启用 pulsar.sink.validateSinkMessageBytes 选项来让链接器使用 Pulsar 提供的 Schema.AUTO_PRODUCE_BYTES() 发送消息。它会在发送字节数组消息时额外进行校验,与 topic 上最新版本的 Schema 进行对比,保证消息内容能正确。
但是,并非所有的 Pulsar 的 Schema 都支持校验字符串,所以在默认情况下我们禁用了此选项。可以按需启用。
自定义序列化器
可以通过继承 PulsarSerializationSchema 接口来实现自定义的序列化逻辑。接口需要返回一个类型为 PulsarMessage 的消息,此类型实例无法被直接创建,连接器提供了构造方法并定义了三种消息类型的构建。
java
// 使用 Pulsar 的 Scheme 来构建消息,常用于你知道 topic 对应的 schema 是什么的时候。我们会检查你提供的 Schema 在 topic 上是否兼容。
PulsarMessage.builder(Schema<M> schema, M message)
...
.build();
// 创建一个消息类型为字节数组的消息。默认情况下不进行 Schema 检查。
PulsarMessage.builder(byte[] bytes)
...
.build();
// 创建一个消息体为空的墓碑消息。墓碑 是一种特殊的消息,并在 Pulsar 中提供了支持。
PulsarMessage.builder()
...
.build();
消息路由策略
在 Pulsar Sink 中,消息路由发生在于分区之间,而非上层 Topic。对于给定 Topic 的情况,路由算法会首先会查询出 Topic 之上所有的分区信息,并在这些分区上实现消息的路由。Pulsar Sink 默认提供 2 种路由策略的实现。
- KeyHashTopicRouter:使用消息的 key 对应的哈希值来取模计算出消息对应的 Topic 分区。使用此路由可以将具有相同 key 的消息发送至同一个 Topic 分区。消息的 key 可以在自定义
PulsarSerializationSchema时,在serialize()方法内使用PulsarMessageBuilder.key(String key)来予以指定。如果消息没有包含 key,此路由策略将从 Topic 分区中随机选择一个发送。 - RoundRobinRouter:轮换使用用户给定的 Topic 分区。消息将会轮替地选取 Topic 分区,当往某个 Topic 分区里写入指定数量的消息后,将会轮换至下一个 Topic 分区。使用
PulsarSinkOptions.PULSAR_BATCHING_MAX_MESSAGES指定向一个 Topic 分区中写入的消息数量。
还可以通过实现 TopicRouter 接口来自定义消息路由策略,请注意 TopicRouter 的实现需要能被序列化。
java
@PublicEvolving
public interface TopicRouter<IN> extends Serializable {
TopicPartition route(IN in, List<TopicPartition> partitions, PulsarSinkContext context);
default void open(SinkConfiguration sinkConfiguration) {
// 默认无操作
}
}
如前文所述,Pulsar 分区的内部被实现为一个无分区的 Topic,一般情况下 Pulsar 客户端会隐藏这个实现,并且提供内置的消息路由策略。Pulsar Sink 并没有使用 Pulsar 客户端提供的路由策略和封装,而是使用了 Pulsar 客户端更底层的 API 自行实现了消息路由逻辑。这样做的主要目的是能够在属于不同 Topic 的分区之间定义更灵活的消息路由策略。
详情请参考 Pulsar 的 partitioned topics 文档。

end
