改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Pulsar调优与实践

ckckck2025年1月11日123 浏览

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_storagemessage_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。

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 文档。

Pulsar 分区示意图

end