- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文基于 Pulsar 官方文档(version-2.3.1 版本,adaptors-kafka.md)完整讲解 Pulsar 提供的 Kafka 兼容封装:如何通过替换 Maven 依赖、调整少量配置项,让现有的 Kafka Java 客户端应用以原有KafkaProducer/KafkaConsumer代码原样指向 Pulsar 服务,并逐条给出 Kafka API 与配置项的兼容性矩阵,以及通过 Kafka properties 直接定制底层 Pulsar 客户端、Producer、Consumer 参数的完整方法。读完后你可以直接复制依赖与示例代码完成迁移,并清楚知道哪些 Kafka API 和配置在封装下不受支持或被忽略。
适用前提:封装模块与当前仓库的关系
该兼容封装以独立 Maven 模块pulsar-client-kafka-compat发布,对外提供两个 artifact:
org.apache.pulsar:pulsar-client-kafka——shaded 版本,重打包了 Kafka 客户端依赖,避免与宿主应用自带的kafka-clients版本冲突;org.apache.pulsar:pulsar-client-kafka-original——未 shaded 版本,供迁移期需要同时引入原生 Kafka 客户端的场景使用。
需要说明的是:当前仓库主干已不再包含pulsar-client-kafka-compat模块(在仓库根目录检索pulsar-client-kafka相关的 pom 与源码均无结果),本文描述的能力对应文档标注的 2.3.1 版本及更早版本线。当前仓库中与 Kafka 的集成主要存在于 pulsar-io/kafka 连接器(作为 Kafka Connector)以及 pulsar-io/kafka-connect-adaptor 中,与本文的“Kafka 客户端封装”是两条不同的集成路径,后文会简要区分。
用 Pulsar Kafka 封装替换 Kafka 客户端依赖
封装的核心设计是“同包名替换”:它复用了org.apache.kafka.clients.producer/org.apache.kafka.clients.consumer等原有包路径下的类名,因此 Java 代码中的 import 和调用无需任何改动,只需在pom.xml中替换依赖。
第一步,删除原来的 Kafka 客户端依赖:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>0.10.2.1</version> </dependency>第二步,引入 Pulsar 的 Kafka 封装(@pulsar:version@是官方文档的占位符,实际使用时替换为对应发布版本号):
<dependency> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-client-kafka</artifactId> <version>@pulsar:version@</version> </dependency>换用新依赖后,原有代码可以不动直接运行,但必须调整两点配置:
- 让 Producer 和 Consumer 指向 Pulsar 服务而不是 Kafka(
bootstrap.servers使用pulsar://协议地址); - 使用特定的 Pulsar 主题(topic)名称,例如
persistent://public/default/my-topic,而不是 Kafka 的短主题名。
迁移期与原生 Kafka 客户端共存
在从 Kafka 向 Pulsar 渐进迁移的过程中,应用很可能一部分流量走原生 Kafka 客户端、另一部分走 Pulsar 封装,两者需要在同一个 JVM 中共存。此时应当改用未 shaded的封装:
<dependency> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-client-kafka-original</artifactId> <version>@pulsar:version@</version> </dependency>因为该依赖不做包名重打包,与宿主应用自带的kafka-clients共用同一套org.apache.kafka.clients类,所以不能再直接new KafkaProducer(...)(那会实例化为原生 Kafka 客户端),而应该使用封装提供的专用入口类构造 Pulsar 侧的客户端:
- Producer 使用
org.apache.kafka.clients.producer.PulsarKafkaProducer,替代KafkaProducer; - Consumer 使用
org.apache.kafka.clients.producer.PulsarKafkaConsumer(按文档原文的类名说明),替代KafkaConsumer。
这样两类客户端在同一个应用里可以各自连接 Kafka 集群和 Pulsar 集群,互不冲突。
Producer 示例
下面是文档中的完整生产者示例。注意注释中强调的:topic 必须是规范的 Pulsar 主题,bootstrap.servers指向 Pulsar 服务。
// Topic needs to be a regular Pulsar topic String topic = "persistent://public/default/my-topic"; Properties props = new Properties(); // Point to a Pulsar service props.put("bootstrap.servers", "pulsar://localhost:6650"); props.put("key.serializer", IntegerSerializer.class.getName()); props.put("value.serializer", StringSerializer.class.getName()); Producer<Integer, String> producer = new KafkaProducer(props); for (int i = 0; i < 10; i++) { producer.send(new ProducerRecord<Integer, String>(topic, i, "hello-" + i)); log.info("Message {} sent successfully", i); } producer.close();几点实践要点:
bootstrap.servers虽然沿用了 Kafka 的键名,但取值是 Pulsar 服务地址(pulsar://localhost:6650),且按消费者配置表中的说明,它需要指向单个Pulsar 服务 URL;- 消息的 partition 参数(示例中的
i)会被映射到 Pulsar 主题的分区路由上; - 序列化器沿用 Kafka 的
key.serializer/value.serializer配置方式,IntegerSerializer、StringSerializer等原生实现可直接使用。
官方文档还给出了更完整的 Producer/Consumer 示例,位于其源码库的pulsar-client-kafka-compat/pulsar-client-kafka-tests/src/test/java/org/apache/pulsar/client/kafka/compat/examples目录下,可对照该目录下的测试工程理解端到端用法(该模块未包含在当前仓库中)。
Consumer 示例
消费者示例展示了订阅、拉取与手动提交位点的完整循环:
String topic = "persistent://public/default/my-topic"; Properties props = new Properties(); // Point to a Pulsar service props.put("bootstrap.servers", "pulsar://localhost:6650"); props.put("group.id", "my-subscription-name"); props.put("enable.auto.commit", "false"); props.put("key.deserializer", IntegerDeserializer.class.getName()); props.put("value.deserializer", StringDeserializer.class.getName()); Consumer<Integer, String> consumer = new KafkaConsumer(props); consumer.subscribe(Arrays.asList(topic)); while (true) { ConsumerRecords<Integer, String> records = consumer.poll(100); records.forEach(record -> { log.info("Received record: {}", record); }); // Commit last offset consumer.commitSync(); }对应到 Pulsar 语义上:
group.id直接映射为 Pulsar 的订阅(subscription)名称,示例中即my-subscription-name;- 关闭自动提交(
enable.auto.commit=false)后,consumer.commitSync()将消费位点提交给 Pulsar 的订阅机制;若开启自动提交,按配置表说明 ack 会立即发送回 broker; - 示例使用
subscribe模式(按主题集合订阅),封装对 rebalance 监听等 Kafka 分区分配语义并未完整支持,详见下文兼容性矩阵。
Kafka API 兼容性矩阵
文档给出的核心结论是:Pulsar 封装支持 Kafka API 的大部分操作。以下表格完整继承自原文档,是评估存量代码能否直接迁移的关键依据。
Producer API
| Producer Method | Supported | Notes |
|---|---|---|
Future<RecordMetadata> send(ProducerRecord<K, V> record) | Yes | |
Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) | Yes | |
void flush() | Yes | |
List<PartitionInfo> partitionsFor(String topic) | No | |
Map<MetricName, ? extends Metric> metrics() | No | |
void close() | Yes | |
void close(long timeout, TimeUnit unit) | Yes |
Producer 配置项
| Config property | Supported | Notes |
|---|---|---|
acks | Ignored | 持久化与 quorum 写在 Pulsar 命名空间级别配置 |
auto.offset.reset | Yes | 未显式设置时默认为latest |
batch.size | Ignored | |
block.on.buffer.full | Yes | 为 true 时阻塞生产者,否则返回错误 |
bootstrap.servers | Yes | |
buffer.memory | Ignored | |
client.id | Ignored | |
compression.type | Yes | 仅支持gzip与lz4,不支持snappy |
connections.max.idle.ms | Yes | 空闲时间上限支持到 2,147,483,647,000 ms(Integer.MAX_VALUE * 1000) |
interceptor.classes | Yes | |
key.serializer | Yes | |
linger.ms | Yes | 控制批量发送消息时的组批提交时间 |
max.block.ms | Ignored | |
max.in.flight.requests.per.connection | Ignored | Pulsar 即使在多个请求同时在途时也能保证顺序 |
max.request.size | Ignored | |
metric.reporters | Ignored | |
metrics.num.samples | Ignored | |
metrics.sample.window.ms | Ignored | |
partitioner.class | Yes | |
receive.buffer.bytes | Ignored | |
reconnect.backoff.ms | Ignored | |
request.timeout.ms | Ignored | |
retries | Ignored | Pulsar 客户端在发送超时到期前以指数退避自动重试 |
send.buffer.bytes | Ignored | |
timeout.ms | Yes | |
value.serializer | Yes |
从矩阵可以看出几个迁移要点:Kafka 中的acks、retries、max.in.flight等一致性/重试参数在 Pulsar 侧由服务端(命名空间级持久化配置)和客户端自身的指数退避重试机制接管,因此被忽略;批量参数中只有linger.ms生效,它控制消息组批的提交窗口。
Consumer API
| Consumer Method | Supported | Notes |
|---|---|---|
Set<TopicPartition> assignment() | No | |
Set<String> subscription() | Yes | |
void subscribe(Collection<String> topics) | Yes | |
void subscribe(Collection<String> topics, ConsumerRebalanceListener callback) | No | |
void assign(Collection<TopicPartition> partitions) | No | |
void subscribe(Pattern pattern, ConsumerRebalanceListener callback) | No | |
void unsubscribe() | Yes | |
ConsumerRecords<K, V> poll(long timeoutMillis) | Yes | |
void commitSync() | Yes | |
void commitSync(Map<TopicPartition, OffsetAndMetadata> offsets) | Yes | |
void commitAsync() | Yes | |
void commitAsync(OffsetCommitCallback callback) | Yes | |
void commitAsync(Map<TopicPartition, OffsetAndMetadata> offsets, OffsetCommitCallback callback) | Yes | |
void seek(TopicPartition partition, long offset) | Yes | |
void seekToBeginning(Collection<TopicPartition> partitions) | Yes | |
void seekToEnd(Collection<TopicPartition> partitions) | Yes | |
long position(TopicPartition partition) | Yes | |
OffsetAndMetadata committed(TopicPartition partition) | Yes | |
Map<MetricName, ? extends Metric> metrics() | No | |
List<PartitionInfo> partitionsFor(String topic) | No | |
Map<String, List<PartitionInfo>> listTopics() | No | |
Set<TopicPartition> paused() | No | |
void pause(Collection<TopicPartition> partitions) | No | |
void resume(Collection<TopicPartition> partitions) | No | |
Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes(Map<TopicPartition, Long> timestampsToSearch) | No | |
Map<TopicPartition, Long> beginningOffsets(Collection<TopicPartition> partitions) | No | |
Map<TopicPartition, Long> endOffsets(Collection<TopicPartition> partitions) | No | |
void close() | Yes | |
void close(long timeout, TimeUnit unit) | Yes | |
void wakeup() | No |
Consumer 配置项
| Config property | Supported | Notes |
|---|---|---|
group.id | Yes | 映射为 Pulsar 订阅名称 |
max.poll.records | Yes | |
max.poll.interval.ms | Ignored | 消息由 broker “推送” |
session.timeout.ms | Ignored | |
heartbeat.interval.ms | Ignored | |
bootstrap.servers | Yes | 需要指向单个 Pulsar 服务 URL |
enable.auto.commit | Yes | |
auto.commit.interval.ms | Ignored | 自动提交时 ack 会立即发送回 broker |
partition.assignment.strategy | Ignored | |
auto.offset.reset | Yes | 仅支持earliest与latest |
fetch.min.bytes | Ignored | |
fetch.max.bytes | Ignored | |
fetch.max.wait.ms | Ignored | |
interceptor.classes | Yes | |
metadata.max.age.ms | Ignored | |
max.partition.fetch.bytes | Ignored | |
send.buffer.bytes | Ignored | |
receive.buffer.bytes | Ignored | |
client.id | Ignored |
通过 Kafka properties 定制 Pulsar 行为
封装允许在 Kafka 的Properties中直接使用pulsar.前缀的配置键,透传到底层 Pulsar 客户端。这是迁移时调整 TLS、认证、超时、组批行为的主要手段。以下三张表完整继承自原文档。
Pulsar 客户端属性
| Config property | Default | Notes |
|---|---|---|
pulsar.authentication.class | 配置认证提供者,例如org.apache.pulsar.client.impl.auth.AuthenticationTls | |
pulsar.authentication.params.map | 表示认证插件参数的 Map | |
pulsar.authentication.params.string | 表示认证插件参数的字符串,例如key1:val1,key2:val2 | |
pulsar.use.tls | false | 启用 TLS 传输加密 |
pulsar.tls.trust.certs.file.path | TLS 信任证书存储的路径 | |
pulsar.tls.allow.insecure.connection | false | 是否接受 broker 的自签名证书 |
pulsar.operation.timeout.ms | 30000 | 通用操作超时时间 |
pulsar.stats.interval.seconds | 60 | Pulsar 客户端库统计打印间隔 |
pulsar.num.io.threads | 1 | Netty IO 线程数 |
pulsar.connections.per.broker | 1 | 到每个 broker 的最大连接数 |
pulsar.use.tcp.nodelay | true | TCP no-delay |
pulsar.concurrent.lookup.requests | 50000 | 最大并发主题查找数 |
pulsar.max.number.rejected.request.per.connection | 50 | 强制关闭连接前的错误阈值 |
典型场景:当集群启用了认证(如 TLS 认证)时,无需修改代码,只要在构造KafkaProducer/KafkaConsumer的Properties中补充:
props.put("pulsar.use.tls", "true"); props.put("pulsar.authentication.class", "org.apache.pulsar.client.impl.auth.AuthenticationTls"); props.put("pulsar.tls.trust.certs.file.path", "/etc/pulsar/certs/ca.pem");Pulsar Producer 属性
| Config property | Default | Notes |
|---|---|---|
pulsar.producer.name | 指定生产者名称 | |
pulsar.producer.initial.sequence.id | 指定该生产者序列号的基线值 | |
pulsar.producer.max.pending.messages | 1000 | 等待 broker 确认的消息队列的最大待发送消息数 |
pulsar.producer.max.pending.messages.across.partitions | 50000 | 跨所有分区的最大待发送消息数 |
pulsar.producer.batching.enabled | true | 控制是否对消息启用自动组批 |
pulsar.producer.batching.max.messages | 1000 | 一个批次中的最大消息数 |
Pulsar Consumer 属性
| Config property | Default | Notes |
|---|---|---|
pulsar.consumer.name | 指定消费者名称 | |
pulsar.consumer.receiver.queue.size | 1000 | 消费者接收队列大小 |
pulsar.consumer.acknowledgments.group.time.millis | 100 | 消费者向 broker 发送确认前的最大组等待时间 |
pulsar.consumer.total.receiver.queue.size.across.partitions | 50000 | 跨分区接收队列的总大小上限 |
pulsar.consumer.subscription.topics.mode | PersistentOnly | 消费者订阅的主题模式 |
与当前仓库中其他 Kafka 集成路径的区分
为避免概念混淆,说明当前仓库中其他与 Kafka 相关的模块与本文封装的关系:
- pulsar-io/kafka:Pulsar 的 Kafka连接器(connector),让 Pulsar 作为消息源/汇与 Kafka Connect 框架对接,运行在 Functions 运行时中,与“Kafka Java 客户端封装”是两个层面的集成;
- pulsar-io/kafka-connect-adaptor:Kafka Connect 适配器,把 Kafka Connect 的 source/sink 任务包装为 Pulsar Functions;
- kafka-connect-avro-converter-shaded:为上述适配器解决 Avro converter 依赖版本冲突而做的 shaded 模块。
也就是说,如果你要“把用了 Kafka 客户端的应用迁移到 Pulsar”,本文的pulsar-client-kafka封装是对路方案;如果你要“让 Pulsar 数据与 Kafka 生态管道互通”,则应关注上述 connector 与 Connect 适配器路径。
小结
- 迁移的核心动作只有三步:
pom.xml中用pulsar-client-kafka替换kafka-clients;bootstrap.servers改为pulsar://地址;topic 改为persistent://tenant/namespace/topic形式,Java 代码保持不变; - 与原生 Kafka 客户端共存时使用
pulsar-client-kafka-original并以PulsarKafkaProducer/PulsarKafkaConsumer显式构造客户端; - 依赖 Kafka 分区分配、rebalance 监听、按时间戳/首尾位点查询、pause/resume 等语义的代码不在支持范围内,迁移前应以本文兼容性矩阵逐条核对;
- TLS、认证、组批、接收队列等深层参数可通过
pulsar.前缀的 properties 直接在 Kafka 配置中透传,无需接触底层客户端 API; - 注意版本适用性:当前仓库主干已移除
pulsar-client-kafka-compat模块,本文内容适用于文档标注的 2.3.1 及包含该封装的历史版本线,使用前请以所选用版本中的该文档与模块为准。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar Kafka 客户端兼容层(pulsar-client-kafka)实战指南:存量 Kafka 应用零改造迁移
Apache Pulsar Kafka 客户端兼容层(pulsar client kafka)实战指南:存量 Kafka 应用零改造迁移 本文以 Apache
消息队列后端流处理Apache Pulsar Kafka 兼容层:零改造迁移 Kafka Java 应用接入 Pulsar 实战指南
Apache Pulsar Kafka 兼容层:零改造迁移 Kafka Java 应用接入 Pulsar 实战指南 本文基于 Apache Pulsar 2.1
消息队列后端流处理Apache Pulsar 的 Kafka 客户端兼容适配器(Kafka Client Wrapper)完整使用指南
Apache Pulsar 的 Kafka 客户端兼容适配器(Kafka Client Wrapper)完整使用指南 本文基于 site2/website ne
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考