ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Pulsar HDFS2 Sink 连接器:将 Topic 消息持久化写入 HDFS 的配置与实战指南

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),仅供参考
返回列表