
SeaTunnel Google Bigtable 连接器实战Source/Sink 配置、sampleRowKeys 并行拆分与数据映射原理解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文围绕 Apache SeaTunnel 的 Google Cloud Bigtable 连接器展开完整覆盖 Source 与 Sink 两端的全部配置项、Schema/类型映射规则、5 类实战任务示例并结合seatunnel-connectors-v2/connector-google-bigtable模块源码深入讲解基于sampleRowKeys的并行拆分机制、分片容错回退、过滤器构建与行编码实现。读完本文你将能够独立配置 Bigtable 表的数据读取与写入任务并理解在 SeaTunnel Zeta 引擎下该连接器的并行读、批量写与故障恢复原理。一、连接器概览与演进记录Google Bigtable 是 Google Cloud 提供的宽列 NoSQL 数据库。SeaTunnel 通过connector-google-bigtable模块提供对该服务的原生接入能力Source 端使用 Bigtable Data v2 Java Client 读取数据Sink 端使用同一客户端写入数据。该连接器的演进记录保存在 changelog/connector-google-bigtable.md 中包含两条核心变更ChangeCommitVersion[Feature][Connector-V2] Add Google Cloud Bigtable Source and Sink connector8e57c04dev[Improve][Connector-V2] Split Bigtable source by sampleRowKeys for parallel reads-dev从变更记录可以看出两条关键信息连接器能力定位该连接器同时提供 Source读取与 Sink写入两个插件官方文档分别见 GoogleBigtable Source 与 GoogleBigtable Sink并行读的实现方向后续改进引入按sampleRowKeys拆分策略使得一个 Bigtable 表可以被切分为多个 tablet 级分片并行扫描——这正是本连接器实现并行度的核心技术下文将结合源码展开。引擎支持范围当前仅支持 SeaTunnel Zeta 引擎见 Source 文档 与 Sink 文档 的 Support Those Engines 章节Flink、Spark 引擎暂不支持。二、Source 连接器从 Bigtable 读取数据2.1 能力特性Source 端能力特性清单对应用户文档中的 Key Featuresbatch有界批式读取stream本身是有界 Source但可通过 STREAMING job mode 配合 checkpoint 运行见后文有界流式扫描示例exactly-once不提供精确一次语义column projection不支持列投影parallelism支持并行且并行拆分正是其特色能力cdc不支持 CDC。2.2 并行读取原理sampleRowKeys 拆分与回退官方文档的核心提示tip描述了完整的并行策略The source is bounded. The enumerator calls BigtablesampleRowKeysto cut the table (or the configuredstart_rowkey/end_rowkeyrange) into tablet-sized splits, then assigns them byhash(splitId) % parallelism. Setenv.parallelism(or the source parallelism) greater than 1 so multiple readers scan different key ranges. If sampling fails, returns no keys, or intersects to an empty range, the connector falls back to a single split so the job can still run. Each scan reads every requested cell for the configured row range and emits one SeaTunnel row per Bigtable row.落到源码层面这一流程由 BigtableSourceSplitEnumerator.java 实现其类注释明确描述了主工作流采样切分调用BigtableClient.sampleRowKeys()Data API 的sampleRowKeysRPC获取表内各 tablet 的近似边界键将整表切分为大小近似相等的半开区间[prev, current)区间求交把每个 tablet 区间与用户配置的start_rowkey/end_rowkey求交集addIntersectedSplit方法取maxStart(rangeStart, userStart)与minEnd(rangeEnd, userEnd)空字符串分别代表表起始/表尾若最后一个采样键非空即 API 未返回表尾哨兵还会追加[lastSample, )区间覆盖表尾数据哈希分配split 通过HashUtils.bucketIndex(split.splitId().hashCode(), parallelism)计算归属 task。当currentParallelism() 1时全部 split 直接分配给唯一 task见assignSplit方法。回退机制buildSplits方法当采样抛出异常、采样结果为空、或与用户区间求交后为空时连接器回退为单个覆盖用户区间的 split保证任务仍能运行不会因切分失败而中断。检查点容错enumerator 在snapshotState中同时持久化assignedSplits与pendingSplits对应 issue #11144 的改进reader 侧则由 BigtableSourceReader.java 在snapshotState中记录正在读取的currentSplit与剩余pendingSplits并通过split.setLastReadRowKey(...)记录进度——恢复时buildQuery会从getResumeStartRowKey()继续扫描实现故障后的断点续读。2.3 Source 配置项详解nametyperequireddefault valueproject_idstringyes-instance_idstringyes-tablestringyes-credentials_pathstringno-rowkey_columnlistno-start_rowkeystringno-end_rowkeystringno-start_timestamplongno-end_timestamplongno-max_versionsintno1scan_row_limitintno-1common-optionsno-各参数说明参数定义与默认值可对照 BigtableSourceOptions.java 与 BigtableBaseOptions.javaproject_id [string]Google Cloud 项目 ID。instance_id [string]Bigtable 实例 ID。table [string]要读取的 Bigtable 表名。credentials_path [string]Google Cloud 服务账号 JSON 密钥文件路径。省略时使用 Application Default CredentialsADC——ADC 在 GCE/GKE 节点、gcloudshell 会话中自动生效也支持通过GOOGLE_APPLICATION_CREDENTIALS环境变量指向服务账号 JSON 文件。rowkey_column [list]可选指定哪些字段接收行键值。不设置时连接器使用名为rowkey的 schema 字段作为行键字段源码中常量ROW_KEY_FIELD rowkey。每个列出的字段按其声明类型独立解码BYTES接收原始行键字节STRING接收 UTF-8 解码视图——因此同一扫描中不同行键字段可以使用不同类型例如一个字段暴露原始字节给下游二进制处理另一个暴露 UTF-8 视图。start_rowkey [string]扫描的包含起始行键。不设置则从表头开始。连接器以 UTF-8 字符串形式传给 Bigtable 客户端仅支持字典序比较二进制行键无法编码为 UTF-8请使用BYTES类型字段。end_rowkey [string]扫描的排他结束行键。不设置则读到表尾。start_timestamp [long]包含起始时间戳过滤器自 epoch 起的微秒数与end_timestamp、max_versions共同控制 Bigtable 为每个列限定符返回哪些 cell 版本。end_timestamp [long]排他结束时间戳过滤器微秒。max_versions [int]每个列限定符最多返回的 cell 版本数。默认1只返回最新版本调大后暴露历史版本但 Source 仍按每 Bigtable 行输出一行处理同一 cell 的旧版本会被折叠进最新返回的 cell 中源码convertRow通过cellMap.putIfAbsent只取每列时间戳最新的第一个 cell。scan_row_limit [int]每个 split最多返回的行数。-1默认表示不限制。当 enumerator 生成多个 split 时作业级上限约为scan_row_limit × split 数而非全表单次上限。可配合start_rowkey/end_rowkey跨多个作业做全表分页扫描。源码中对应Query.limit(...)仅当 0时设置。common optionsSource 插件通用参数详见 Source Common Options。2.4 Schema 映射规则SeaTunnel schema 字段名必须遵循familyName:qualifier模式例如cf:name、stats:age行键字段由rowkey_column控制未配置时使用特殊字段名rowkey。Schema field nameMapped Bigtable cellrowkeyRow keycf:nameColumn familycf, qualifiernamestats:ageColumn familystats, qualifierage读取类型注意事项Source 为每个family:qualifier字段读取最新返回的 cell用start_timestamp、end_timestamp、max_versions控制 Bigtable 扫描过滤器。SeaTunnel 字段类型必须与 Bigtable 中存储的字节一致——例如本连接器写入的数值类型是二进制大端big-endian值而STRING、DATE、TIME、TIMESTAMP、DECIMAL是 UTF-8 文本。底层解码逻辑见 BigtableDeserializationFormat.javaINT/BIGINT/FLOAT/DOUBLE用ByteBuffer.wrap(...)按大端还原BOOLEAN取首字节非 0 即 true日期时间类按对应DateTimeFormatter模式解析文本。2.5 Source 任务示例示例一使用 ADC 读取全部行env { parallelism 1 job.mode BATCH } source { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events schema { fields { rowkey BYTES cf:type STRING cf:ts BIGINT } } } }示例二使用服务账号扫描行键区间env { parallelism 1 job.mode BATCH } source { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events credentials_path /secrets/sa-key.json start_rowkey 2024-01-01# end_rowkey 2024-02-01# max_versions 1 schema { fields { rowkey STRING cf:type STRING cf:data STRING } } } }示例三自定义行键字段名env { parallelism 1 job.mode BATCH } source { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events rowkey_column [event_id] schema { fields { event_id STRING cf:type STRING cf:data STRING } } } }示例四带 cell 版本过滤的有界流式扫描使用STREAMINGjob mode 时扫描在启用 checkpoint 的情况下运行但本身仍是一次有界读取。组合start_timestamp、end_timestamp、max_versions可以限定 Bigtable 返回哪些 cell 版本env { parallelism 1 job.mode STREAMING checkpoint.interval 60000 } source { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events start_timestamp 1704067200000000 end_timestamp 1735689600000000 max_versions 3 scan_row_limit 500000 schema { fields { rowkey STRING cf:type STRING cf:data STRING cf:ts BIGINT } } } }注示例中的时间戳1704067200000000、1735689600000000均为微秒精度分别对应 2024-01-01 与 2024-12-31 附近。2.6 Source 底层扫描构造BigtableSourceReader.buildQuery揭示了读取过程的三个关键细节源码行键区间Query.range(startKey, endKey)其中startKey优先取 checkpoint 恢复用的getResumeStartRowKey()实现断点续读结束键取 split 的endRowKey行数上限scan_row_limit 0时调用query.limit(...)过滤器链max_versions 0时叠加Filters.FILTERS.limit().cellsPerColumn(maxVersions)时间戳过滤器按起止都存在 / 仅有 start / 仅有 end三种情况分别构造timestamp().range().of(startTs, endTs)、startClosed(startTs)、endOpen(endTs)多个过滤器通过Filters.FILTERS.chain()串联。读取时通过bigtableClient.getDataClient().readRows(query).forEach(...)逐行流式消费每收集一行才持有一次 checkpoint 锁执行output.collect(...)避免把整个结果集缓冲在内存中。三、Sink 连接器向 Bigtable 写入数据3.1 能力特性Sink 端能力特性清单batch批式写入support multiple table write支持多表写入通过multi_table_sink_replica参数exactly-once不提供精确一次cdc不支持 CDC 语义timer flush不依赖定时刷新写入由批量大小与 checkpoint 触发。3.2 Sink 配置项详解nametyperequireddefault valueproject_idstringyes-instance_idstringyes-tablestringyes-rowkey_columnlistyes-column_familyconfigyes-credentials_pathstringno-rowkey_delimiterstringnoversion_columnstringno-null_modestringnoskipbatch_mutation_sizeintno100schema_save_modeenumnoRECREATE_SCHEMAdata_save_modeenumnoAPPEND_DATAmulti_table_sink_replicaintno1common-optionsno-各参数说明定义见 BigtableSinkOptions.javaproject_id [string]Google Cloud 项目 ID示例my-gcp-project。instance_id [string]Bigtable 实例 ID示例my-bigtable-instance。table [string]要写入的 Bigtable 表名示例my-table。连接器不会自动建表运行作业前须手动创建表及全部所需列族column family。rowkey_column [list]用于拼接 Bigtable 行键的列名示例[id]或[tenant, id]。多列时用rowkey_delimiter连接。单行键列时null 或空值会使作业以WRITE_FAILED失败多行键列时任一非末尾列出现 null 会静默变成拼接行键中的空段由rowkey_delimiter连接仅当整个拼接结果为空时作业才失败见BigtableSinkWriter.buildRowKey单列走fieldToByteStringBYTES 类型取原始字节、其余按 UTF-8多列用String.join(delimiter, parts)。column_family [config]列名到列族名的映射。用all_columns作为 key 可为所有未映射列设置默认列族column_family { name info age stats }或把所有列放入同一列族column_family { all_columns cf }未出现在映射中的字段名回退到all_columns指定的列族若未配置all_columns则回退到默认列族cf源码常量DEFAULT_FAMILY cf见BigtableSinkWriter.resolveFamily。credentials_path [string]Google Cloud 服务账号 JSON 密钥文件路径不设置则使用 Application Default CredentialsADC——在 GCE/GKE 上自动生效或通过GOOGLE_APPLICATION_CREDENTIALS环境变量指定。rowkey_delimiter [string]拼接多个行键列值使用的分隔符默认空串即无分隔符。version_column [string]列名其BIGINT值用作 Bigtable cell 时间戳自 epoch 起的微秒数。不设置时使用当前系统时间源码实现为System.currentTimeMillis() * 1000L因 Bigtable 时间戳单位为微秒。null_mode [string]null 字段值处理方式支持skip默认与emptyskip该 cell 不出现在 mutation 中empty向 cell 写入空字节数组mutation.setCell(family, qualifier, timestamp, ByteString.EMPTY)。batch_mutation_size [int]攒够多少条行 mutation 后向 Bigtable 发送一次 BulkMutation默认100。调大可提升吞吐但会提高单 task 内存占用。写入路径见BigtableSinkWriter.writebuffer.size() batchMutationSize时触发flush()prepareCommitcheckpoint 提交与close时也会强制 flush。schema_save_mode [enum]schema 保存模式当前仅支持RECREATE_SCHEMA。连接器不创建 Bigtable 表或列族作业开始前需手动建好目标表与全部列族。data_save_mode [enum]数据保存模式当前仅支持APPEND_DATA。DROP_DATA与ERROR_WHEN_DATA_EXISTS尚未实现如需干净的目标表请在作业前自行清空或重建。multi_table_sink_replica [int]多表写入的 sink 副本数详见 Sink Common Options。该参数增加单个 sink 实例内的并行写入副本数目标 Bigtable 表由table选项固定不会按上游表推导。common optionsSink 插件通用参数详见 Sink Common Options。3.3 数据类型映射所有 SeaTunnel 类型均受支持写入 Bigtable 时的存储格式如下对应BigtableSinkWriter.convertToByteString的实现分支SeaTunnel typeStorage format in BigtableTINYINT1-byte binarySMALLINT2-byte big-endian binaryINT4-byte big-endian binaryBIGINT8-byte big-endian binaryFLOAT4-byte IEEE 754 big-endianDOUBLE8-byte IEEE 754 big-endianBOOLEAN1-byte (1 true, 0 false)BYTESRaw bytesSTRINGUTF-8 textDECIMALUTF-8 plain stringDATEUTF-8yyyy-MM-ddTIMEUTF-8HH:mm:ssTIMESTAMPUTF-8yyyy-MM-dd HH:mm:ss重要语义提示Bigtable 没有关系型列概念。Sink 将每个非行键字段写为一个大单元格cell目标列族由column_family决定Bigtable qualifier 即 SeaTunnel 字段名。Sink 将每一行上游数据视为一次无条件 cell 变更因此UPDATE/DELETE行类型不会被解释为 CDC 操作而是在相同的(row key, column family, qualifier)三元组下覆盖之前的 cell。3.4 Sink 任务示例示例一ADC 基础写入env { parallelism 1 job.mode BATCH } sink { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events rowkey_column [event_id] column_family { all_columns cf } } }示例二服务账号密钥文件 复合行键env { parallelism 1 job.mode BATCH } sink { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events credentials_path /secrets/sa-key.json rowkey_column [tenant_id, event_id] rowkey_delimiter # column_family { all_columns data } batch_mutation_size 500 } }示例三多列族映射env { parallelism 1 job.mode BATCH } sink { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table user_profile rowkey_column [user_id] column_family { name identity email identity age stats last_login stats } } }示例四版本列 空值写空env { parallelism 1 job.mode BATCH } sink { GoogleBigtable { project_id my-gcp-project instance_id my-bigtable-instance table events rowkey_column [tenant_id, event_id] rowkey_delimiter # version_column event_ts null_mode empty column_family { all_columns data event_type meta } } }示例五流式写入 checkpoint 刷盘流式模式下writer 在每个 checkpoint 时刷新内存中的 mutation 缓冲。batch_mutation_size仍控制 task 内缓冲大小checkpoint 频率只影响已缓冲 mutation 何时发送给 Bigtable。env { parallelism 2 job.mode STREAMING checkpoint.interval 30000 } source { FakeSource { row.num 1000 schema { fields { tenant_id string event_id string event_ts bigint event_type string payload string } } plugin_output events_stream } } sink { GoogleBigtable { plugin_input events_stream project_id my-gcp-project instance_id my-bigtable-instance table events credentials_path /secrets/sa-key.json rowkey_column [tenant_id, event_id] rowkey_delimiter # version_column event_ts column_family { all_columns data event_type meta } batch_mutation_size 200 } }3.5 Sink 写入实现要点从 BigtableSinkWriter.java 源码可以梳理出写入侧的几条实现事实mutation 构建每个上游行转换为一个RowKeyMutation行键空时报WRITE_FAILEDRow key cannot be empty. Check rowkey_column configuration.随后排除行键列与版本列对剩余每个字段mutation.setCell(family, qualifier, timestamp, valueBytes)批量发送buffer 达到batch_mutation_size即调用bigtableClient.bulkMutate(...)flush时先buffer.clear()再发送避免bulkMutate抛异常时重复发送编码一致性写入侧的编码数值大端、布尔单字节、日期时间文本化与 Source 侧BigtableDeserializationFormat的解码严格对称因此Bigtable 连接器写、Bigtable 连接器读的闭环任务天然类型自洽多表写入支持writer 实现了SupportMultiTableSinkWriter接口与multi_table_sink_replica参数配合在单个 sink 实例内提供并行写入副本。四、使用前提与限制清单综合官方文档与源码实现使用该连接器前需明确以下前提与限制仅支持 SeaTunnel Zeta 引擎不支持 Flink / Spark starterBigtable 表与列族需手动创建Source 读取表已存在是前提Sink 不会建表且schema_save_mode仅支持RECREATE_SCHEMA、data_save_mode仅支持APPEND_DATA清库需自行完成行键比较仅支持字典序Source 的start_rowkey/end_rowkey以 UTF-8 字符串传递二进制行键请通过BYTES类型字段在 schema 侧表达读取为非精确一次语义Source 不提供 exactly-oncemax_versions 1时旧版本 cell 会被折叠进最新值写入无 CDC 语义UPDATE/DELETE行类型会按普通覆盖写入处理凭据策略优先推荐把服务账号 JSON 路径写入credentials_path生产环境也可依赖 ADCGCE/GKE 元数据或GOOGLE_APPLICATION_CREDENTIALS环境变量并行度设置想要多 reader 并行扫描需将env.parallelism或 Source 并行度设为大于 1由 enumerator 按hash(splitId) % parallelism分配分片单并行度时所有分片交给唯一 task。五、总结SeaTunnel 的 Google Bigtable 连接器是一对功能完整的 Source/Sink 插件Source 端以sampleRowKeys为核心实现了 tablet 级并行拆分与单分片回退容错支持行键区间、时间戳区间、版本数与行数上限等细粒度扫描控制Sink 端以 bulk mutation 批量写入为骨架支持复合行键、列族映射、版本列、空值策略与 checkpoint 刷盘。通过本文的配置示例与源码级解析你可以直接在本仓库的 Source 文档、Sink 文档 与 连接器源码 中进一步追溯实现细节。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考