尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

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

发布时间:2026/9/25 3:39:09

资讯中心
01
ARTICLE

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

Apache Pulsar Kafka 客户端兼容封装(pulsar-client-kafka-compat):让 Kafka 应用零改动迁移到 Pulsar
消息队列后端流处理【免费下载链接】pulsarApache 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发布对外提供两个 artifactorg.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 groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version0.10.2.1/version /dependency第二步引入 Pulsar 的 Kafka 封装pulsar:version是官方文档的占位符实际使用时替换为对应发布版本号dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-client-kafka/artifactId versionpulsar:version/version /dependency换用新依赖后原有代码可以不动直接运行但必须调整两点配置让 Producer 和 Consumer 指向 Pulsar 服务而不是 Kafkabootstrap.servers使用pulsar://协议地址使用特定的 Pulsar 主题topic名称例如persistent://public/default/my-topic而不是 Kafka 的短主题名。迁移期与原生 Kafka 客户端共存在从 Kafka 向 Pulsar 渐进迁移的过程中应用很可能一部分流量走原生 Kafka 客户端、另一部分走 Pulsar 封装两者需要在同一个 JVM 中共存。此时应当改用未 shaded的封装dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-client-kafka-original/artifactId versionpulsar:version/version /dependency因为该依赖不做包名重打包与宿主应用自带的kafka-clients共用同一套org.apache.kafka.clients类所以不能再直接new KafkaProducer(...)那会实例化为原生 Kafka 客户端而应该使用封装提供的专用入口类构造 Pulsar 侧的客户端Producer 使用org.apache.kafka.clients.producer.PulsarKafkaProducer替代KafkaProducerConsumer 使用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()); ProducerInteger, String producer new KafkaProducer(props); for (int i 0; i 10; i) { producer.send(new ProducerRecordInteger, 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()); ConsumerInteger, String consumer new KafkaConsumer(props); consumer.subscribe(Arrays.asList(topic)); while (true) { ConsumerRecordsInteger, 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.commitfalse后consumer.commitSync()将消费位点提交给 Pulsar 的订阅机制若开启自动提交按配置表说明 ack 会立即发送回 broker示例使用subscribe模式按主题集合订阅封装对 rebalance 监听等 Kafka 分区分配语义并未完整支持详见下文兼容性矩阵。Kafka API 兼容性矩阵文档给出的核心结论是Pulsar 封装支持 Kafka API 的大部分操作。以下表格完整继承自原文档是评估存量代码能否直接迁移的关键依据。Producer APIProducer MethodSupportedNotesFutureRecordMetadata send(ProducerRecordK, V record)YesFutureRecordMetadata send(ProducerRecordK, V record, Callback callback)Yesvoid flush()YesListPartitionInfo partitionsFor(String topic)NoMapMetricName, ? extends Metric metrics()Novoid close()Yesvoid close(long timeout, TimeUnit unit)YesProducer 配置项Config propertySupportedNotesacksIgnored持久化与 quorum 写在 Pulsar 命名空间级别配置auto.offset.resetYes未显式设置时默认为latestbatch.sizeIgnoredblock.on.buffer.fullYes为 true 时阻塞生产者否则返回错误bootstrap.serversYesbuffer.memoryIgnoredclient.idIgnoredcompression.typeYes仅支持gzip与lz4不支持snappyconnections.max.idle.msYes空闲时间上限支持到 2,147,483,647,000 msInteger.MAX_VALUE * 1000interceptor.classesYeskey.serializerYeslinger.msYes控制批量发送消息时的组批提交时间max.block.msIgnoredmax.in.flight.requests.per.connectionIgnoredPulsar 即使在多个请求同时在途时也能保证顺序max.request.sizeIgnoredmetric.reportersIgnoredmetrics.num.samplesIgnoredmetrics.sample.window.msIgnoredpartitioner.classYesreceive.buffer.bytesIgnoredreconnect.backoff.msIgnoredrequest.timeout.msIgnoredretriesIgnoredPulsar 客户端在发送超时到期前以指数退避自动重试send.buffer.bytesIgnoredtimeout.msYesvalue.serializerYes从矩阵可以看出几个迁移要点Kafka 中的acks、retries、max.in.flight等一致性/重试参数在 Pulsar 侧由服务端命名空间级持久化配置和客户端自身的指数退避重试机制接管因此被忽略批量参数中只有linger.ms生效它控制消息组批的提交窗口。Consumer APIConsumer MethodSupportedNotesSetTopicPartition assignment()NoSetString subscription()Yesvoid subscribe(CollectionString topics)Yesvoid subscribe(CollectionString topics, ConsumerRebalanceListener callback)Novoid assign(CollectionTopicPartition partitions)Novoid subscribe(Pattern pattern, ConsumerRebalanceListener callback)Novoid unsubscribe()YesConsumerRecordsK, V poll(long timeoutMillis)Yesvoid commitSync()Yesvoid commitSync(MapTopicPartition, OffsetAndMetadata offsets)Yesvoid commitAsync()Yesvoid commitAsync(OffsetCommitCallback callback)Yesvoid commitAsync(MapTopicPartition, OffsetAndMetadata offsets, OffsetCommitCallback callback)Yesvoid seek(TopicPartition partition, long offset)Yesvoid seekToBeginning(CollectionTopicPartition partitions)Yesvoid seekToEnd(CollectionTopicPartition partitions)Yeslong position(TopicPartition partition)YesOffsetAndMetadata committed(TopicPartition partition)YesMapMetricName, ? extends Metric metrics()NoListPartitionInfo partitionsFor(String topic)NoMapString, ListPartitionInfo listTopics()NoSetTopicPartition paused()Novoid pause(CollectionTopicPartition partitions)Novoid resume(CollectionTopicPartition partitions)NoMapTopicPartition, OffsetAndTimestamp offsetsForTimes(MapTopicPartition, Long timestampsToSearch)NoMapTopicPartition, Long beginningOffsets(CollectionTopicPartition partitions)NoMapTopicPartition, Long endOffsets(CollectionTopicPartition partitions)Novoid close()Yesvoid close(long timeout, TimeUnit unit)Yesvoid wakeup()NoConsumer 配置项Config propertySupportedNotesgroup.idYes映射为 Pulsar 订阅名称max.poll.recordsYesmax.poll.interval.msIgnored消息由 broker “推送”session.timeout.msIgnoredheartbeat.interval.msIgnoredbootstrap.serversYes需要指向单个 Pulsar 服务 URLenable.auto.commitYesauto.commit.interval.msIgnored自动提交时 ack 会立即发送回 brokerpartition.assignment.strategyIgnoredauto.offset.resetYes仅支持earliest与latestfetch.min.bytesIgnoredfetch.max.bytesIgnoredfetch.max.wait.msIgnoredinterceptor.classesYesmetadata.max.age.msIgnoredmax.partition.fetch.bytesIgnoredsend.buffer.bytesIgnoredreceive.buffer.bytesIgnoredclient.idIgnored通过 Kafka properties 定制 Pulsar 行为封装允许在 Kafka 的Properties中直接使用pulsar.前缀的配置键透传到底层 Pulsar 客户端。这是迁移时调整 TLS、认证、超时、组批行为的主要手段。以下三张表完整继承自原文档。Pulsar 客户端属性Config propertyDefaultNotespulsar.authentication.class配置认证提供者例如org.apache.pulsar.client.impl.auth.AuthenticationTlspulsar.authentication.params.map表示认证插件参数的 Mappulsar.authentication.params.string表示认证插件参数的字符串例如key1:val1,key2:val2pulsar.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-delaypulsar.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 propertyDefaultNotespulsar.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 propertyDefaultNotespulsar.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/kafkaPulsar 的 Kafka连接器connector让 Pulsar 作为消息源/汇与 Kafka Connect 框架对接运行在 Functions 运行时中与“Kafka Java 客户端封装”是两个层面的集成pulsar-io/kafka-connect-adaptorKafka Connect 适配器把 Kafka Connect 的 source/sink 任务包装为 Pulsar Functionskafka-connect-avro-converter-shaded为上述适配器解决 Avro converter 依赖版本冲突而做的 shaded 模块。也就是说如果你要“把用了 Kafka 客户端的应用迁移到 Pulsar”本文的pulsar-client-kafka封装是对路方案如果你要“让 Pulsar 数据与 Kafka 生态管道互通”则应关注上述 connector 与 Connect 适配器路径。小结迁移的核心动作只有三步pom.xml中用pulsar-client-kafka替换kafka-clientsbootstrap.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 及包含该封装的历史版本线使用前请以所选用版本中的该文档与模块为准。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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消息队列后端流处理上一篇10个顶级SwiftUI开源iOS应用推荐来自gh_mirrors/ex/example-ios-apps的精选项目下一篇10个Starlark核心特性详解确定性、密封性、并行执行创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。