Apache Pulsar HDFS2 Sink 连接器:将 Topic 消息持久化写入 HDFS 的配置与实战指南
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载HDFS2 Sink 是 Apache Pulsar 内置的 IO 连接器之一负责从 Pulsar Topic 拉取消息并将消息以文件形式持久化到 Hadoop HDFSHadoop 2.x 生态中是Pulsar 消息 → 数据湖/离线分析存储链路上最常见的落库手段之一。本文以 site2/website-next/docs/io-hdfs2-sink.md 为骨架结合 pulsar-io/hdfs2 模块的源码与测试完整讲解其全部配置属性、校验规则、JSON/YAML 配置模板以及底层写入与确认ack机制帮助你在实际集群中正确部署并调优该连接器。HDFS2 Sink 连接器是什么HDFS2 Sink 是一个 Pulsar IO Sink 组件核心职责一句话概括从 Pulsar Topic 拉取消息并将其持久化到 HDFS 文件中。它适合如下场景将实时消息流持续归档到 HDFS供 Hive、Spark、Presto 等离线/批量计算引擎消费将 Pulsar 作为统一事件总线Sink 作为数据出口之一把数据搬运到 HDFS需要与 Hadoop 2.x 体系含 Kerberos 认证打通的存算分离架构。从仓库结构看该连接器位于 pulsar-io/hdfs2其模块名为pulsar-io-hdfs2构建时依赖hadoop-client2.8.5见 pulsar-io/hdfs2/pom.xml因此它面向 Hadoop 2.x 生态同仓库中的 pulsar-io/hdfs3 则面向 Hadoop 3.x。数据写入的核心机制源码级解读HDFS2 Sink 不是简单地把每条消息直接写盘而是通过缓冲队列 后台同步线程实现批量落盘与确认。整条链路由以下类协作完成HdfsSinkConfig.java负责加载并校验配置HdfsAbstractSink.javaSink 主体负责连接 HDFS、创建文件、启动同步线程HdfsSyncThread.java后台线程周期调用hsync()刷盘并 ack 已落盘记录HdfsAbstractTextFileSink.java及seq包下的序列文件实现负责实际写入格式AbstractHdfsConnector.javaHDFS 连接与 Kerberos 认证的底层封装。启动流程open()在 HdfsAbstractSink.java#L57-L69 中open()依次执行通过HdfsSinkConfig.load(config)解析配置支持 Map 与 YAML 文件两种入口调用hdfsSinkConfig.validate()做参数合法性校验用maxPendingRecords作为容量创建LinkedBlockingQueueRecordV即未确认记录队列若配置了subdirectoryPattern则编译对应的DateTimeFormatter连接 HDFSconnectToHdfs()、创建文件写入器createWriter()、启动同步线程launchSyncThread()。文件路径与文件命名目标文件路径在 HdfsAbstractSink.java#L100-L118 的getPath()中生成规则为{directory}/{subdirectoryPattern 格式化的当前时间}/ {filenamePrefix}-{System.currentTimeMillis()}{扩展名}其中扩展名的优先级是fileExtension配置优先若未配置fileExtension但设置了compression则使用压缩编解码器的默认扩展名如 GZIP 的.gz。这也解释了文档中filenamePrefix 使topicA产生名为topicA-...的文件的描述文件名由filenamePrefix - 时间戳拼接而成。写入、刷盘与确认写入以文本格式为例HdfsAbstractTextFileSink.java#L56-L69 将每条记录record.getValue().toString()写入OutputStreamWriter若配置了separator非默认\u0000空字符则在记录后追加分隔符写入成功后记录放入未确认队列写入失败则调用record.fail()。同步与确认HdfsSyncThread.java#L46-L78 的后台线程每隔syncInterval毫秒执行一次先对 HDFS 输出流调用hsync()强制刷盘再把队列中的记录逐一ack()。当syncInterval为 0 时线程会尽快循环刷盘close()时调用halt()做最后一次刷盘并确认全部待确认记录。背压语义未确认队列容量即maxPendingRecords。当队列满时put()会阻塞从而对上游形成天然背压保证内存中的未确认记录数不超过阈值。这正是文档中设置为 1 时每条记录先落盘再 ack最多一次内存缓冲、设置为较大值时允许批量缓冲后统一刷盘的底层原因整体遵循at-least-once至少一次投递语义。连接与安全认证AbstractHdfsConnector.java#L64-L97 的resetHDFSResources()负责初始化 HadoopConfiguration与FileSystem按逗号切分hdfsConfigResources逐个addResource加载配置资源需能被类路径找到否则抛出IOException关闭FileSystem缓存fs.scheme.impl.disable.cachetrue避免重配置后无法生效若集群启用了 KerberosSecurityUtil.isSecurityEnabled则以kerberosUserPrincipal keytab登录并执行doAs获取 FileSystem否则走simple认证。配置属性全解析HDFS2 Sink 的配置由 HdfsSinkConfig.javaSink 特有属性与 AbstractHdfsConfig.javaHDFS 通用属性共同承载全部属性如下表名称类型必填默认值说明hdfsConfigResourcesString是None包含 Hadoop 文件系统配置的一个文件或逗号分隔的文件列表。示例core-site.xmlhdfs-site.xmldirectoryString是NoneHDFS 中读取或写入文件的目录。encodingString否None文件的字符编码。示例UTF-8ASCIIcompressionCompression否NoneHDFS 上压缩/解压文件所使用的压缩编解码器可选值如下BZIP2DEFLATEGZIPLZ4SNAPPYkerberosUserPrincipalString否None用于认证的 Kerberos 用户主体principal。keytabString否None用于认证的 Kerberos keytab 文件的完整路径。filenamePrefixString是compression为None时NoneHDFS 目录内创建文件的前缀。示例取值为topicA时将生成名为topicA-...的文件。fileExtensionString是None写入 HDFS 文件的扩展名。示例.txt.seqseparatorchar否None文本文件中用于分隔记录的字符。若未设置所有记录的内容将首尾相连拼接为一个连续字节数组。syncIntervallong否0调用 flush 将数据写入 HDFS 磁盘的间隔毫秒。maxPendingRecordsint否Integer.MAX_VALUEack 之前允许在内存中持有的最大记录数。设为 1 时每条记录在 ack 前都会先落盘设为大值时允许先缓冲多条记录再统一刷盘。subdirectoryPatternString否None与 Sink 创建时间关联的子目录。该模式是directory子目录的格式化模式。模式语法见 DateTimeFormatter该链接为 Oracle 官方 JDK 文档供查阅模式语法。关键参数深度说明hdfsConfigResources与directory二者是 AbstractHdfsConfig.java#L66-L75 中强校验的必填项任一缺失都会抛出Required property not set.。hdfsConfigResources指向的文件会被config.addResource()加载到 HadoopConfiguration中必须位于连接器的类路径上。compression与fileExtension的联动HdfsSinkConfig.java#L94-L99 的校验逻辑是当fileExtension为空且compression也为空时才报错即只要设置了压缩fileExtension可省略将由CompressionCodec.getDefaultExtension()补全如 GZIP →.gz。同时 AbstractHdfsConnector.java#L182-L191 通过CompressionCodecFactory.getCodecByName()按名称解析编解码器枚举定义见 Compression.javaBZIP2, DEFLATE, GZIP, LZ4, SNAPPY。kerberosUserPrincipal与keytab二者必须成对出现——只配置其一都会抛出Values for both kerberosUserPrincipal keytab are required.见 AbstractHdfsConfig.java#L71-L74。syncInterval与maxPendingRecordssyncInterval不能为负数maxPendingRecords必须为正整数。它们共同决定吞吐优先还是延迟/可靠性优先追求低延迟逐条落盘可设maxPendingRecords1追求高吞吐可调大该值并配合合理的syncInterval。subdirectoryPattern模式会先经 HdfsSinkConfig.java#L109-L115 用固定时间LocalDateTime.of(2020, 1, 1, 12, 0)做格式合法性校验非法模式会在启动阶段直接报错运行时则由 HdfsAbstractSink.java#L110-L112 用LocalDateTime.now()生成实际子目录。常见的按天归档写法为yyyy-MM-dd。encoding若未配置则回退到 JVM 默认字符集见 AbstractHdfsConnector.java#L177-L180。跨集群部署时建议显式指定如UTF-8避免因运行环境不同导致乱码。separator未配置时默认值为 Java 空字符\u0000写入时不做分隔见 HdfsAbstractTextFileSink.java#L61-L63配置后每条记录后追加该字符方便下游按分隔符切分。配置示例使用 HDFS2 Sink 连接器前需先通过以下任一方式创建配置文件以下两例均完整复现自官方文档可直接套用。JSON 示例{ configs: { hdfsConfigResources: core-site.xml, directory: /foo/bar, filenamePrefix: prefix, fileExtension: .log, compression: SNAPPY, subdirectoryPattern: yyyy-MM-dd } }YAML 示例configs: hdfsConfigResources: core-site.xml directory: /foo/bar filenamePrefix: prefix fileExtension: .log compression: SNAPPY subdirectoryPattern: yyyy-MM-dd以上示例对应的解析结果hdfsConfigResourcescore-site.xml、directory/foo/bar、filenamePrefixprefix、compressionSNAPPY、subdirectoryPatternyyyy-MM-dd在测试类 HdfsSinkConfigTests.java#L39-L48 与 #L51-L66 中被逐一断言验证可直接作为配置正确性的参照。快速上手部署 HDFS2 Sink前置条件可访问的 Hadoop 2.x 集群本模块依赖hadoop-client2.8.5将core-site.xml、hdfs-site.xml等 Hadoop 配置资源放入连接器可访问的类路径可放置于 NAR 包内或通过extraDependenciesDir挂载因为底层通过config.addResource()按资源名加载若 HDFS 开启 Kerberos需准备 principal 与 keytab 文件路径目标 HDFS 目录已存在或运行账户有创建权限文件不存在时连接器会调用fs.create(path)存在时追加写入见 HdfsAbstractSink.java#L95。创建 Sink使用 Pulsar Admin CLI 的 sinks 命令即可创建archive 指向构建出的 HDFS2 Sink NAR 包可用mvn package从 pulsar-io/hdfs2 构建bin/pulsar-admin sinks create \ --tenant public \ --namespace default \ --name hdfs2-sink \ --sink-type hdfs2 \ --archive /path/to/pulsar-io-hdfs2.nar \ --inputs input-topic \ --sink-config-file /path/to/sink-config.yaml创建后Sink 实例会以 Pulsar Function 运行时的方式运行消费input-topic上的消息并按前述流程写入 HDFS。暂停或删除可分别使用bin/pulsar-admin sinks pause与bin/pulsar-admin sinks delete。具体参数说明可查阅仓库 pulsar-client-tools 模块中 sinks 相关命令实现。常见配置组合与调优建议归档场景追求吞吐maxPendingRecords调大如 10000syncInterval设为几百到几千毫秒批量刷盘换取更高吞吐文件扩展名建议配合compression如 SNAPPY压缩落盘节省 HDFS 空间。严格可靠性场景maxPendingRecords1每条消息先落盘再 acksyncInterval保持较小值。代价是吞吐显著下降仅建议在低流量关键链路使用。按天分目录归档设置subdirectoryPattern: yyyy-MM-ddHDFS 上会生成/foo/bar/2026-09-23/prefix-timestamp.log这样的层级结构配合directory统一管理生命周期。多条消息混写不设置separator时记录首尾相接下游解析需按固定长度设置separator如\n后可按行切分更利于文本类下游工具直接处理。测试与验证仓库在 pulsar-io/hdfs2/src/test 下提供了完整的配置与端到端测试可用于验证你的配置理解HdfsSinkConfigTests.java覆盖 YAML/Map 加载、必填项缺失、非法压缩编解码器、负 syncInterval、非正 maxPendingRecords、Kerberos 参数不成对等全部校验分支AbstractHdfsSinkTest.java 及text/seq包下的 HdfsStringSinkTests.java、HdfsTextSinkTests.java、HdfsSequentialSinkTests.java 则验证了实际写入行为。总结HDFS2 Sink 是 Apache Pulsar 连接 HDFS 2.x 生态的标准出口它以缓冲队列 后台hsync刷盘 批量 ack的机制在吞吐与可靠性之间提供了maxPendingRecords、syncInterval两个核心旋钮通过compression、fileExtension、subdirectoryPattern等参数可灵活控制落盘文件的格式、压缩与目录布局同时借助 Hadoop 原生配置资源与 Kerberos 支持无缝融入企业级安全集群。理解 HdfsSinkConfig.java 的校验逻辑与 HdfsAbstractSink.java 的写入流程是正确配置与排障的关键起点。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Kafka Sink Connector 实战指南将 Pulsar Topic 消息桥接到 KafkaApache Pulsar Kafka Sink Connector 实战指南将 Pulsar Topic 消息桥接到 Kafka Kafka Sink Co消息队列后端流处理Apache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr CollectionApache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr Collection Solr si消息队列后端流处理Apache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 RedisApache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 Redis 本篇技术指南以 Apache Puls消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

程序佬独立游戏角色动画破局:三大引擎方案与程序化动画实战

程序佬独立游戏角色动画破局:三大引擎方案与程序化动画实战

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

2026/9/24 2:49:09 阅读更多 →
手机CPU虚焊维修三大方案:热风枪、BGA返修台与吹风机法全解析

手机CPU虚焊维修三大方案:热风枪、BGA返修台与吹风机法全解析

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

2026/9/24 2:48:09 阅读更多 →
嵌入式软件静态测试(十四)——IEC 61508 / IEC 62304:功能安全与医疗器械嵌入式软件的静态测试流程

嵌入式软件静态测试(十四)——IEC 61508 / IEC 62304:功能安全与医疗器械嵌入式软件的静态测试流程

❄️ 我的个人专栏: 《智能软件工程AI4SE》 《嵌入式面试总结》 《嵌入式处理器架构解析》 《嵌入式与虚拟化》 《嵌入式软件测试》 🌟 Simplicity is the ultimate sophistication摘要:本文围绕 IEC 61508 与 IEC 62304 两项标准&#xff0…

2026/9/24 2:48:09 阅读更多 →

最新新闻

黄金微针按次报价怎样核对包含项和变更差价

黄金微针按次报价怎样核对包含项和变更差价

黄金微针写着“按次收费”,并不自动说明一次包括哪些内容。到了比较报价或调整方案时,真正影响支出的,是服务范围怎样变化、原付款有多少可以用于新方案,以及哪些款项仍在单独处理中。先统一口径,再算差额,…

2026/9/24 4:50:28 阅读更多 →
国产安全MCU LKT6830C开发实战:硬件加密与防篡改设计

国产安全MCU LKT6830C开发实战:硬件加密与防篡改设计

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

2026/9/24 4:50:28 阅读更多 →
Storm 安全加固:Kerberos 认证、ACL 权限与多租户隔离

Storm 安全加固:Kerberos 认证、ACL 权限与多租户隔离

Storm 安全加固:Kerberos 认证、ACL 权限与多租户隔离Apache Storm 作为分布式实时计算框架,广泛应用于实时数据处理场景。随着企业级应用的需求增长,Storm 平台的安全性也日益重要。本文将详细介绍 Storm 安全加固的三大核心机制&#xff1a…

2026/9/24 4:50:28 阅读更多 →
Flutter鸿蒙化适配:screen_protector防截屏插件ArkTS实现指南

Flutter鸿蒙化适配:screen_protector防截屏插件ArkTS实现指南

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

2026/9/24 4:50:28 阅读更多 →
RK3506 AMP双系统实战:Linux+FreeRTOS核间通信与实时性优化

RK3506 AMP双系统实战:Linux+FreeRTOS核间通信与实时性优化

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

2026/9/24 4:50:28 阅读更多 →
CodeBurn 发布验收 Agent 执行手册:从候选 SHA 到 release-ready 的可复现审计契约

CodeBurn 发布验收 Agent 执行手册:从候选 SHA 到 release-ready 的可复现审计契约

【免费下载链接】codeburn Free, local tool to track AI coding token usage and cost across 37 tools and agents (Claude Code, Cursor, Codex, Gemini and more), by model, project, and task. npx codeburn 项目地址: https://gitcode.com/gh_mirrors/co/cod…

2026/9/24 4:49:27 阅读更多 →

日新闻

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为…

2026/9/24 0:00:19 阅读更多 →
单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

简介:一份基于单细胞RNA测序数据的细胞类型注释算法研究Python毕业设计源码,针对计算机相关专业正在做毕设或需要项目实战的学习者,可用于课程设计与期末大作业。项目代码完整、经导师指导评审通过,可直接运行,覆盖数据…

2026/9/24 0:00:19 阅读更多 →
C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

第一次在项目里被反射卡住,是在一个老旧的WinForms模块里:几十个类依赖PropertyChanged通知,运行时反射读属性、发通知,每次启动慢半拍不说,一上.NET Native/AOT裁剪模式几乎全面崩盘。后来我把这段逻辑全部改成C#源生…

2026/9/24 0:00:19 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/23 4:55:02 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/23 4:49:06 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/23 9:53:41 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/23 9:53:40 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/23 9:53:40 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/23 9:53:40 阅读更多 →