
SeaTunnel HBase Source 连接器实战指南批量扫描、行键/时间范围读取与 Kerberos 配置【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 内置的HbaseSource 连接器用于从 Apache HBase 表批量读取数据支持普通全表扫描、行键范围扫描、时间戳范围扫描、二进制行键、自定义命名空间与基于 Region 的并行分片读取。读完本文你将掌握 HBase Source 的全部配置项语义与默认值、三种典型读取场景的 HOCON 配置写法以及 Kerberos 安全集群下的接入方式并能结合源码理解连接器的并行分片与数据反序列化原理。引擎支持与连接器定位HbaseSource 连接器在三种引擎下均可运行Spark / Flink / SeaTunnel Zeta从 HbaseSource.java 的源码实现看它实现了SeaTunnelSource接口并返回Boundedness.BOUNDED也就是说该连接器以批处理模式运行读取的是扫描范围在任务启动时刻的快照状态并非 CDC 源扫描启动后发生的新写入不会进入结果任务读取完毕即结束不会持续等待新数据。这一点在HbaseSourceReader的pollNext中也有体现当分片消费完毕且noMoreSplit为 true 时会调用context.signalNoMoreElement()通知结束而不是继续阻塞等待。主要特性一览特性支持情况批处理✅ 支持流处理❌ 不支持精确一次Exactly Once❌ 不支持模式投影Column Projection✅ 支持并行度Parallelism✅ 支持用户自定义拆分❌ 不支持其中模式投影与并行度两项在源码中有直接对应HbaseSource同时实现了SupportParallelism与SupportColumnProjection接口。并行度的实现细节见下文并行分片原理一节。选项总览HBase Source 支持的全部选项如下默认值均与 HbaseSourceOptions.java 中的定义一致名称类型必填默认值描述zookeeper_quorumstring是-HBase ZooKeeper 地址例如hadoop001:2181,hadoop002:2181。tablestring是-要扫描的 HBase 表。自定义 namespace 请使用namespace:table。schemaconfig是-SeaTunnel 表结构。行键写作rowkey普通单元格写作列簇:列名。hbase_extra_configconfig否-额外的 HBase 或 Hadoop 客户端配置。cachingint否-1每次 RPC 获取的行数。-1表示使用 HBase 客户端默认值。batchint否-1每次 RPC 返回的最大单元格数量。-1表示使用 HBase 客户端默认值。cache_blocksboolean否false扫描结果是否写入 HBase block cache。is_binary_rowkeyboolean否false是否把行键字段按二进制字节处理。start_rowkeystring否-范围扫描的起始行键。end_rowkeystring否-范围扫描的结束行键。start_row_inclusiveboolean否true扫描结果是否包含start_rowkey。end_row_inclusiveboolean否false扫描结果是否包含end_rowkey。start_timestamplong否-时间范围扫描的起始时间戳包含该时间。end_timestamplong否-时间范围扫描的结束时间戳不包含该时间。common-options-否-Source 插件通用参数例如plugin_output。必填项zookeeper_quorum、table的必填校验定义在 HbaseSourceFactory.java 的optionRule()中schema则由CatalogTableUtil.buildWithConfig在创建 Source 时解析未配置会直接报错。核心参数详解zookeeper_quorum [string]必填HBase 的 ZooKeeper 集群主机列表多个节点用逗号分隔例如hadoop001:2181,hadoop002:2181,hadoop003:2181连接器在构建连接时见 HbaseClient.java 的getHbaseConnection会执行hbaseConfiguration.set(hbase.zookeeper.quorum, ...)把它写入 HBase 客户端配置的核心项。table [string]必填要读取的表名例如seatunnel。支持两种写法默认命名空间直接写表名如seatunnel_test。解析逻辑见 HbaseParameters.java 的buildWithSourceConfig当table中不含冒号时namespace 会被赋值为 HBase 默认命名空间default自定义命名空间使用namespace:table形式如ns1:seatunnel_test连接器会按冒号拆分出 namespace 与表名。schema [config]必填HBase 使用字节数组进行存储因此必须为表中的每一列显式声明数据类型。规则如下行键列使用固定名称rowkeyHBase 单元格使用列簇:列名形式例如info:name字段类型遵循 SeaTunnel 类型系统可声明为string、boolean、double、bigint、bytes、int、float、decimal、date、time、timestamp等。关于列名的校验与反序列化见下文数据读取与类型反序列化一节的源码解析。更多模式声明细节可参考 模式声明指南。hbase_extra_config [config]可选用于注入额外的 HBase / Hadoop 客户端配置例如安全认证、RPC 保护级别等。该配置会以键值对形式逐项写入org.apache.hadoop.conf.ConfigurationhbaseExtraConfig.forEach(hbaseConfiguration::set)因此在 Kerberos 场景下HBase 与 Hadoop 的安全相关配置都应放在这里。caching可选默认 -1设置扫描过程中一次从服务器端获取的行数。较大的 caching 可以减少客户端与服务器之间的 RPC 往返次数从而提高扫描效率但如果单次获取数据量过大也会增加客户端内存压力。默认值-1表示沿用 HBase 客户端自身的默认值。在HbaseClient.buildScan中对应scan.setCaching(hbaseParameters.getCaching())。batch可选默认 -1设置扫描时每次 RPC 返回的最大单元格数量。对于单行包含大量列的宽表尤其有用可以避免一次性返回过多数据节省客户端内存并改善吞吐。默认值-1表示沿用 HBase 客户端默认值对应scan.setBatch(...)。cache_blocks可选默认 false设置扫描过程中是否缓存数据块。默认情况下 HBase 会把扫描到的数据块写入 BlockCache便于后续访问命中但大规模全表扫描会污染BlockCache把热点数据挤出缓存。SeaTunnel 中默认值为false即扫描时不缓存数据块从而减少内存占用、保护缓存命中率。对应scan.setCacheBlocks(...)。is_binary_rowkey可选默认 falseHBase 的行键既可以是文本字符串也可以是任意二进制字节。SeaTunnel 中默认按文本字符串处理is_binary_rowkey false此时Bytes.toBytes(rowKey)将字符串按 UTF-8 编码当设置为true时HBaseUtil.convertRowKey会改用Bytes.toBytesBinary(rowKey)把十六进制/转义形式的输入解析为原始字节见 HBaseUtil.java。若行键为二进制还应在 schema 中将行键列声明为bytes。start_rowkey / end_rowkey可选行键范围扫描的起止行键。设置后连接器只读取该区间内的数据。若start_rowkey大于end_rowkeyHBaseUtil.validateRowKeyRange会抛出IllegalArgumentException拒绝任务启动。这两个参数在 HbaseSourceSplitEnumerator.java 中还会参与 Region 分片的裁剪详见并行分片原理。start_row_inclusive可选默认 true设置扫描范围是否包含起始行。默认true包含。注意大多数情况下应保持默认值仅当有特定需求需要排除起始行时才改为false。end_row_inclusive可选默认 false设置扫描范围是否包含结束行。默认false不包含遵循 HBase 标准的左闭右开区间约定[start, end)。仅当需要在扫描结果中包含结束行时才改为true。重要提示多分片并行读取时的边界组合问题使用多个分片并行读取时这两个参数的组合对数据完整性至关重要默认配置start_row_inclusivetrue, end_row_inclusivefalse推荐配置可确保跨分片不丢数据、不重复数据。每个分片遵循[start, end)左闭右开约定都设置为 falsestart_row_inclusivefalse, end_row_inclusivefalse可能导致数据丢失因为边界行会被所有分片排除在外都设置为 truestart_row_inclusivetrue, end_row_inclusivetrue可能导致数据重复因为边界行会被相邻的多个分片重复包含。start_timestamp / end_timestamp可选时间范围扫描的起止时间戳单位为毫秒epoch millis时间范围同样遵循[start, end)左闭右开约定。start_timestamp包含边界end_timestamp不包含边界。只设置start_timestamp最大值视为无限上界源码中max Long.MAX_VALUE只设置end_timestamp最小值视为无限下界源码中min 0Lstart_timestamp必须大于等于 0end_timestamp必须大于 0若两者同时配置必须满足start_timestamp end_timestampstart_timestamp end_timestamp将导致空扫描。这些约束在 HbaseSourceFactory.java 的optionRule()中通过Conditions.greaterOrEqual、greaterThan、lessThanField声明在 HbaseClient.java 的applyTimeRange中还会二次校验并映射为scan.setTimeRange(min, max)当start_rowkey/end_rowkey与start_timestamp/end_timestamp同时配置时行键范围与时间范围会被同时应用最终返回两者的交集。常用选项common-optionsSource 插件通用参数例如plugin_output、result_table_name等具体参考 Source 常用选项。配置示例按行键和时间范围读取以下配置同时使用行键范围[B, C)与时间范围读取seatunnel_test表中rowkey从B到C、且写入时间落在[1700000000000, 1700003600000)毫秒内的单元格source { Hbase { zookeeper_quorum hadoop001:2181,hadoop002:2181,hadoop003:2181 table seatunnel_test caching 1000 batch 100 cache_blocks false is_binary_rowkey false start_rowkey B end_rowkey C start_timestamp 1700000000000 end_timestamp 1700003600000 schema { columns [ { name rowkey type string }, { name columnFamily1:column1 type boolean }, { name columnFamily1:column2 type double }, { name columnFamily2:column1 type bigint } ] } } }注意该配置会将行键范围与时间范围同时作用于每个分片的 Scan 对象最终结果集是两者交集若想忽略行键条件只按时间读取去掉start_rowkey/end_rowkey即可。读取命名空间下的表表名中使用namespace:table形式即可读取自定义命名空间下的表source { Hbase { zookeeper_quorum hbase_e2e:2181 table ns1:seatunnel_test schema { columns [ { name rowkey, type string }, { name info:name, type string } ] } } }读取二进制行键当 HBase 表的行键是二进制数据而非文本时source { Hbase { zookeeper_quorum hbase_e2e:2181 table binary_rowkey_table is_binary_rowkey true caching 500 batch 100 schema { columns [ { name rowkey, type bytes }, { name info:name, type string }, { name info:score, type double } ] } } }当is_binary_rowkey true时请务必在schema中把行键列声明为bytes并在后续的 transform 节点中自行解码例如转换为十六进制字符串或按业务规则还原。若此时误将行键声明为string反序列化阶段会按 UTF-8 直接转换字节得到的将是乱码字符串。Kerberos 安全集群配置在启用 Kerberos 的 HBase / Hadoop 集群上使用本连接器需要注意以下几点见 Hbase.md 原文备注connector-hbase不会解析krb5_path/kerberos_principal/kerberos_keytab_path这类连接器级参数需要在运行环境中提前完成 Kerberos 登录并保证krb5.conf可被 JVM 访问例如执行kinit -kt keytab principal完成 Ticket 获取或通过 JVM 参数-Djava.security.krb5.conf/path/to/krb5.conf指定配置文件将 HBase / Hadoop 的安全配置写入hbase_extra_config它们会通过hbaseConfiguration.set(...)注入 HBase 客户端。完整示例source { Hbase { zookeeper_quorum zk1:2181,zk2:2181,zk3:2181 table source_table caching 1000 batch 200 cache_blocks false is_binary_rowkey false # HBase安全配置 hbase_extra_config { hbase.security.authentication kerberos hadoop.security.authentication kerberos hbase.master.kerberos.principal hbase/_HOSTREALM hbase.regionserver.kerberos.principal hbase/_HOSTREALM hbase.rpc.protection authentication hbase.zookeeper.useSasl false } schema { columns [ { name rowkey, type string }, { name info:name, type string }, { name info:score, type string } ] } } }其中REALM请替换为实际 Kerberos Realm如EXAMPLE.COM_HOST由客户端按目标主机名自动解析。若集群还开启了 Ranger 或 ACL还需确保运行 SeaTunnel 的操作系统用户对目标表具备读写权限。源码原理并行分片与数据读取Region 级并行分片HbaseSourceSplitEnumerator是理解连接器并行能力的关键。其getTableSplits()方法的工作流程如下见 HbaseSourceSplitEnumerator.java通过RegionLocator.getStartKeys()/getEndKeys()获取目标表全部 Region 的起止行键用HBaseUtil.convertRowKey把用户配置的start_rowkey/end_rowkey转成字节并做合法性校验逐 Region 与用户行键范围求交集裁剪出每个分片的[startRow, endRow)封装为HbaseSourceSplit分片按 splitId 的 hashCode 分配到各并行 subtaskgetSplitOwner使用HashUtils.bucketIndex实现数据在多个 Reader 间的均衡若表不存在或拿不到 Region 信息会抛出HbaseConnectorException错误码TABLE_QUERY_EXCEPTION。每个HbaseSourceSplit都带独立的[startRow, endRow)字节范围见 HbaseSourceSplit.java这正解释了前文边界参数组合影响数据完整性的结论分片边界如何截断直接由start_row_inclusive/end_row_inclusive决定。扫描构造与时间范围每个分片在HbaseSourceReader.pollNext中被消费通过HbaseClient.scan构建 Scan 并获取ResultScanner。buildScan的组装逻辑见 HbaseClient.javascan.withStartRow(split.getStartRow(), isStartRowInclusive)与scan.withStopRow(split.getEndRow(), isEndRowInclusive)应用分片行键范围与开闭约定scan.setCacheBlocks / setCaching / setBatch应用三个性能相关参数scan.addColumn(family, qualifier)按 schema 中声明的列簇:列名精确投影所需列这就是模式投影特性的底层实现scan.setTimeRange(min, max)应用时间范围条件仅在配置了时间戳时生效。列名校验与类型反序列化HbaseSourceReader构造时会对 schema 中的非 rowkey 列做格式校验必须符合列簇:列名且冒号分割后恰好两段否则抛出IllegalArgumentException提示 Invalid column names, it should be [ColumnFamily:Column] format。读取时rowkey列取result.getRow()普通列按列簇:列名拆分后取result.getValue(familyBytes, qualifierBytes)并通过并发 Map 缓存拆分结果以避免重复计算。字节数组到 SeaTunnel 类型的转换由 HBaseDeserializationFormat.java 完成各类型转换规则如下SeaTunnel 类型转换方式TINYINT直接取首字节SMALLINT两字节按大端组合INTBytes.toIntBOOLEANBytes.toBooleanBIGINTBytes.toLongFLOAT / DOUBLEBytes.toFloat/Bytes.toDoubleDECIMAL优先按字符串解析失败时回退到 float 转换BYTES原样返回字节数组DATE / TIME / TIMESTAMP按yyyy-MM-dd/HH:mm:ss/yyyy-MM-dd HH:mm:ss解析字符串STRINGBytes.toStringUTF-8其他抛出UNSUPPORTED_DATA_TYPE异常这也解释了为什么读取二进制行键时 schema 必须声明bytes只有bytes类型才会原样保留字节内容便于后续 transform 自定义解码。变更日志连接器的历史变更记录位于 connector-hbase 变更日志。建议在升级连接器版本前查阅该文件确认参数行为是否有调整例如边界开闭约定、时间戳校验规则等。常见问题排查建议连接失败 /Build Hbase connection failed检查zookeeper_quorum是否可达、hbase_extra_config是否缺少安全认证配置表不存在 / 无 Region 信息连接器会分别在枚举分片阶段报TABLE_QUERY_EXCEPTION请确认表名含 namespace与运行用户权限行键范围报错start_rowkey大于end_rowkey时任务直接失败请检查两者字典序与is_binary_rowkey是否与实际存储一致并行读取后数据缺失或重复优先保持start_row_inclusivetrue、end_row_inclusivefalse的默认组合不要随意同时关闭或同时开启边界包含二进制行键乱码确认is_binary_rowkeytrue且 schema 行键列为bytes并在 transform 中完成解码。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考