news 2026/9/25 3:39:08

Apache Pulsar Kafka 客户端兼容封装(pulsar-client-kafka-compat):让 Kafka 应用零改动迁移到 Pulsar

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar Kafka 客户端兼容封装(pulsar-client-kafka-compat):让 Kafka 应用零改动迁移到 Pulsar
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

本文基于 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>

换用新依赖后,原有代码可以不动直接运行,但必须调整两点配置:

  1. 让 Producer 和 Consumer 指向 Pulsar 服务而不是 Kafka(bootstrap.servers使用pulsar://协议地址);
  2. 使用特定的 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 MethodSupportedNotes
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 propertySupportedNotes
acksIgnored持久化与 quorum 写在 Pulsar 命名空间级别配置
auto.offset.resetYes未显式设置时默认为latest
batch.sizeIgnored
block.on.buffer.fullYes为 true 时阻塞生产者,否则返回错误
bootstrap.serversYes
buffer.memoryIgnored
client.idIgnored
compression.typeYes仅支持gzip与lz4,不支持snappy
connections.max.idle.msYes空闲时间上限支持到 2,147,483,647,000 ms(Integer.MAX_VALUE * 1000)
interceptor.classesYes
key.serializerYes
linger.msYes控制批量发送消息时的组批提交时间
max.block.msIgnored
max.in.flight.requests.per.connectionIgnoredPulsar 即使在多个请求同时在途时也能保证顺序
max.request.sizeIgnored
metric.reportersIgnored
metrics.num.samplesIgnored
metrics.sample.window.msIgnored
partitioner.classYes
receive.buffer.bytesIgnored
reconnect.backoff.msIgnored
request.timeout.msIgnored
retriesIgnoredPulsar 客户端在发送超时到期前以指数退避自动重试
send.buffer.bytesIgnored
timeout.msYes
value.serializerYes

从矩阵可以看出几个迁移要点:Kafka 中的acks、retries、max.in.flight等一致性/重试参数在 Pulsar 侧由服务端(命名空间级持久化配置)和客户端自身的指数退避重试机制接管,因此被忽略;批量参数中只有linger.ms生效,它控制消息组批的提交窗口。

Consumer API

Consumer MethodSupportedNotes
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 propertySupportedNotes
group.idYes映射为 Pulsar 订阅名称
max.poll.recordsYes
max.poll.interval.msIgnored消息由 broker “推送”
session.timeout.msIgnored
heartbeat.interval.msIgnored
bootstrap.serversYes需要指向单个 Pulsar 服务 URL
enable.auto.commitYes
auto.commit.interval.msIgnored自动提交时 ack 会立即发送回 broker
partition.assignment.strategyIgnored
auto.offset.resetYes仅支持earliest与latest
fetch.min.bytesIgnored
fetch.max.bytesIgnored
fetch.max.wait.msIgnored
interceptor.classesYes
metadata.max.age.msIgnored
max.partition.fetch.bytesIgnored
send.buffer.bytesIgnored
receive.buffer.bytesIgnored
client.idIgnored

通过 Kafka properties 定制 Pulsar 行为

封装允许在 Kafka 的Properties中直接使用pulsar.前缀的配置键,透传到底层 Pulsar 客户端。这是迁移时调整 TLS、认证、超时、组批行为的主要手段。以下三张表完整继承自原文档。

Pulsar 客户端属性

Config propertyDefaultNotes
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.tlsfalse启用 TLS 传输加密
pulsar.tls.trust.certs.file.pathTLS 信任证书存储的路径
pulsar.tls.allow.insecure.connectionfalse是否接受 broker 的自签名证书
pulsar.operation.timeout.ms30000通用操作超时时间
pulsar.stats.interval.seconds60Pulsar 客户端库统计打印间隔
pulsar.num.io.threads1Netty IO 线程数
pulsar.connections.per.broker1到每个 broker 的最大连接数
pulsar.use.tcp.nodelaytrueTCP no-delay
pulsar.concurrent.lookup.requests50000最大并发主题查找数
pulsar.max.number.rejected.request.per.connection50强制关闭连接前的错误阈值

典型场景:当集群启用了认证(如 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 propertyDefaultNotes
pulsar.producer.name指定生产者名称
pulsar.producer.initial.sequence.id指定该生产者序列号的基线值
pulsar.producer.max.pending.messages1000等待 broker 确认的消息队列的最大待发送消息数
pulsar.producer.max.pending.messages.across.partitions50000跨所有分区的最大待发送消息数
pulsar.producer.batching.enabledtrue控制是否对消息启用自动组批
pulsar.producer.batching.max.messages1000一个批次中的最大消息数

Pulsar Consumer 属性

Config propertyDefaultNotes
pulsar.consumer.name指定消费者名称
pulsar.consumer.receiver.queue.size1000消费者接收队列大小
pulsar.consumer.acknowledgments.group.time.millis100消费者向 broker 发送确认前的最大组等待时间
pulsar.consumer.total.receiver.queue.size.across.partitions50000跨分区接收队列的总大小上限
pulsar.consumer.subscription.topics.modePersistentOnly消费者订阅的主题模式

与当前仓库中其他 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

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载
上一篇:10个顶级SwiftUI开源iOS应用推荐:来自gh_mirrors/ex/example-ios-apps的精选项目
下一篇:10个Starlark核心特性详解:确定性、密封性、并行执行

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 3:38:59

SpringBoot+Vue二手交易系统:从业务建模到部署上线全解析

最近我把一套基于 SpringBoot Vue 的二手物品交易管理系统重新翻了出来&#xff0c;项目代号 bootpf&#xff0c;代码包名统一叫 com.bootpf。这套系统从用户注册、商品发布、浏览搜索、购物车、下订单&#xff0c;到后台的商品审核、用户管理和数据统计&#xff0c;基本把二手…

作者头像 李华
网站建设 2026/9/25 3:37:33

STM32入门第0集:从芯片认知到环境搭建,新手避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/25 3:37:08

华为昇腾推理引擎开源:边缘AI部署与性能优化实战

1. 昇腾推理引擎开源这件事&#xff0c;到底在解决什么问题第一次在昇腾社区看到推理引擎开源的消息时&#xff0c;我正蹲在一个边缘计算项目上折腾模型部署。当时手里的活儿是把一个视觉检测模型塞进一台功耗受限的工控机里&#xff0c;芯片用的是昇腾310P3。那会儿最头疼的不…

作者头像 李华