新闻详情

新闻详情

首页 / 资讯中心 / 详情

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

发布时间:2026/9/25 3:39:09来源:尧图网络
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),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

Swift Package Manager 的 `swift package` 命令完全指南:创建、编辑、检查与维护 Swift 包 2026/9/25 5:21:40

Swift Package Manager 的 `swift package` 命令完全指南:创建、编辑、检查与维护 Swift 包

开发工具构建工具 【免费下载链接】swift-package-manager The Package Manager for the Swift Programming Language 项目地址: https://gitcode.com/gh_mirrors/sw/swift-package-manager 点击查看 免费下载 swift package 是 Swift Package Manager(…

阅读更多 →
Amazon SNS 实战指南(AWS SDK for Java 2.x):从 Hello SNS 到 FIFO 主题与 SNS 到 SQS 扇出 2026/9/25 5:21:40

Amazon SNS 实战指南(AWS SDK for Java 2.x):从 Hello SNS 到 FIFO 主题与 SNS 到 SQS 扇出

示例工程教程后端 【免费下载链接】aws-doc-sdk-examples Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below. 项目地…

阅读更多 →
Cyrus SASL 2.1.21 编译部署与排错指南:从源码到认证链路 2026/9/25 5:21:40

Cyrus SASL 2.1.21 编译部署与排错指南:从源码到认证链路

简介:这是一份 Cyrus SASL 2.1.21 开源认证库的源码压缩包,面向邮件服务器管理员、安全运维人员以及有二次开发需求的嵌入式开发者。它主要服务于 SMTP、IMAP、POP3 等协议场景,提供多种可插拔的认证机制,是 Postfix 等邮件传输代…

阅读更多 →
node-glob 完全指南:Bash 语义的 glob 匹配、API 家族全景与源码级原理剖析 2026/9/25 5:21:40

node-glob 完全指南:Bash 语义的 glob 匹配、API 家族全景与源码级原理剖析

开发工具 【免费下载链接】node-glob glob functionality for node.js 项目地址: https://gitcode.com/gh_mirrors/no/node-glob 点击查看 免费下载 本篇技术指南以当前仓库 README.md 为主体,系统讲解 glob(node-glob)这个 Node…

阅读更多 →
The Concise TypeScript Book 精读:TypeScript 索引签名(Index Signatures)的三种键类型与 JavaScript 键转换机制 2026/9/25 5:21:34

The Concise TypeScript Book 精读:TypeScript 索引签名(Index Signatures)的三种键类型与 JavaScript 键转换机制

文档教程 【免费下载链接】typescript-book The Concise TypeScript Book: A Concise Guide to Effective Development in TypeScript. Free and Open Source. 项目地址: https://gitcode.com/gh_mirrors/typ/typescript-book 点击查看 免费下载 导读 索引签名&am…

阅读更多 →
nnU-Net v2 安装与初始化配置完全指南:从零搭建首个 nnU-Net 运行环境 2026/9/25 5:21:28

nnU-Net v2 安装与初始化配置完全指南:从零搭建首个 nnU-Net 运行环境

人工智能深度学习计算机视觉医疗健康 【免费下载链接】nnUNet 项目地址: https://gitcode.com/gh_mirrors/nn/nnUNet 点击查看 免费下载 导读 本文是 nnU-Net v2(即 nnunetv2 包)的首次运行环境搭建指南,系统覆盖 PyTorch 前置安…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉