新闻详情

新闻详情

首页 / 资讯中心 / 详情

SeaTunnel GoogleFirestore Sink 连接器实战指南:将 SeaTunnel 数据行写入 Google Cloud Firestore

发布时间:2026/9/17 3:06:51来源:尧图网络
SeaTunnel GoogleFirestore Sink 连接器实战指南:将 SeaTunnel 数据行写入 Google Cloud Firestore
SeaTunnel GoogleFirestore Sink 连接器实战指南将 SeaTunnel 数据行写入 Google Cloud Firestore【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文围绕 SeaTunnel 的connector-google-firestoreSink 连接器展开系统讲解如何将 SeaTunnel 作业中的每一条数据行SeaTunnel Row转换为 Google Cloud Firestore 文档并写入指定 Collection。文章从配置选项、认证方式、数据类型映射到批/流两种作业形态的完整配置示例并结合仓库源码FirestoreSinkWriter、FirestoreSinkFactory、DefaultSeaTunnelRowSerializer剖析其底层实现原理与使用边界。读完本文你将能够独立完成 GoogleFirestore Sink 的认证配置、字段类型映射校验与批/流作业编写并理解其每行一次 Firestore add 调用写入模型的限制。一、连接器概述GoogleFirestore Sink 连接器的作用是把 SeaTunnel 作业产生的数据行写入 Google Cloud 的 Firestore 数据库集合Collection。其核心工作方式是一行一文档每一条 SeaTunnel Row 被序列化为一个 Firestore 文档Document自动生成文档 ID连接器通过调用 Firestore 客户端的add(...)方法写入文档因此文档 ID 由 Firestore 自动生成属于追加写入模式而不是按用户指定的文档 ID 进行更新凭据灵活既可以在配置中显式传入 Base64 编码的 Service Account JSON也可以依赖运行环境的 Google Application Default CredentialsADC不管理索引连接器不会创建或管理 Firestore 索引运行涉及索引的查询前需要先在 Google Cloud 控制台创建所需索引。支持的引擎引擎支持情况Spark✅Flink✅SeaTunnel Zeta✅功能特性支持矩阵特性支持exactly-once精确一次❌CDC❌batch批处理✅stream流处理✅多表写入support multiple table write❌定时冲刷timer flush❌关于上述特性的详细定义可参考 Connector V2 特性说明。依赖获取数据源依赖GoogleFirestoreorg.apache.seatunnel:connector-google-firestore依赖可通过install-plugin.sh脚本安装或从 Maven Central 仓库下载。在插件映射文件 plugin-mapping.properties 中该连接器的注册信息为seatunnel.sink.GoogleFirestore connector-google-firestore对应仓库中的模块为 connector-google-firestore其 POM 中声明的 Firestore 客户端版本为3.7.10见 pom.xml。二、选项Options详解连接器支持的全部配置项如下名称类型是否必填默认值说明project_idstring是-拥有 Firestore 数据库的 Google Cloud 项目 ID不可为空collectionstring是-要写入的 Firestore Collection 名称不可为空credentialsstring否-Base64 编码的 Google Cloud Service Account JSON若配置则不可为空common-options-否-Sink 公共选项详见 Sink Common Optionsproject_id [string]必填项。指定拥有 Firestore 数据库的 Google Cloud 项目 ID值不能为空白字符串。collection [string]必填项。指定要写入的 Firestore Collection 名称值不能为空白字符串。一个 Sink 块只写入一个 Collection——若上游输入来自多个表需要为每个 Firestore Collection 配置一个独立的 Sink 块。credentials [string]可选项。Base64 编码的 Google Cloud Service Account JSON。若配置值不能为空白字符串。若不配置该选项连接器使用 Google Application Default CredentialsADC请确保环境变量GOOGLE_APPLICATION_CREDENTIALS指向 Service Account JSON 文件或运行环境本身已提供默认凭据生成 Base64 值的命令# Linux / 通用 base64 -w 0 service-account.json# macOS base64 service-account.json | tr -d \n⚠️ 不要将原始的 Service Account JSON 直接填入credentials必须先做 Base64 编码。common optionsSink 插件的公共参数请参阅 Sink Common Options。配置校验的源码级印证选项的必填/可选规则在 FirestoreSinkFactory.java 中通过OptionRule定义return OptionRule.builder() .required(PROJECT_ID, notBlank(PROJECT_ID)) .required(COLLECTION, notBlank(COLLECTION)) .optional(CREDENTIALS, notBlank(CREDENTIALS)) .build();即project_id与collection为必填且notBlank不能为空白credentials为可选但一旦配置同样要求非空白。三个选项在 FirestoreSinkOptions.java 中定义均为stringType()且noDefaultValue()。对应的单元测试 FirestoreFactoryTest.java 覆盖了以下校验场景缺少project_id或collection时抛出OptionValidationExceptionproject_id、collection、credentials传入空串、空格、\t、\n、\r等空白值时均被拒绝传入未定义的unknown_option会被拒绝validateUnknownKeys校验。三、认证方式与客户端初始化原理连接器的认证逻辑位于 FirestoreSinkWriter.java 的构造函数中GoogleCredentials credentials; if (parameters.getCredentials() ! null) { byte[] bytes Base64.getDecoder().decode(parameters.getCredentials()); credentials GoogleCredentials.fromStream(new ByteArrayInputStream(bytes)); } else { credentials GoogleCredentials.getApplicationDefault(); } FirestoreOptions firestoreOptions FirestoreOptions.getDefaultInstance() .toBuilder() .setProjectId(parameters.getProjectId()) .setCredentials(credentials) .build(); this.firestore firestoreOptions.getService(); this.collectionReference firestore.collection(parameters.getCollection()); this.serializer new DefaultSeaTunnelRowSerializer(seaTunnelRowType);从源码可以看出完整的初始化链路凭据解析若配置了credentials先进行 Base64 解码再通过GoogleCredentials.fromStream(...)从字节流中加载 Service Account否则回退到GoogleCredentials.getApplicationDefault()读取环境中的 ADC构建客户端基于FirestoreOptions.getDefaultInstance()设置projectId与credentials随后调用getService()创建 Firestore 客户端绑定集合通过firestore.collection(parameters.getCollection())获取目标 CollectionReference准备序列化器用上游表的SeaTunnelRowType初始化DefaultSeaTunnelRowSerializer后续每行数据都由它转换为 Firestore 的MapString, Object文档字段。此外参数对象 FirestoreParameters.java 的buildWithConfig会把作业配置中的project_id、collection、credentials提取为可序列化的参数对象供 Writer 使用。四、写入模型逐行 add 与无缓冲语义FirestoreSinkWriter的核心写入方法如下Override public void write(SeaTunnelRow seaTunnelRow) throws IOException { collectionReference.add(serializer.serialize(seaTunnelRow)); }当前实现中write()对每一行调用一次 Firestore 客户端的add(...)方法既不缓冲也不批量合并行。这意味着写入请求是逐条同步发起的每个add返回一个ApiFuture代码并未显式等待或聚合没有内存中的写入缓冲区因此不存在在 checkpoint 边界进行 flush 的机制checkpoint 完成并不代表此前写入的行已经全部到达 Firestore——这是在使用 checkpoint/容错能力时务必注意的语义差别该连接器在当前实现下不支持 exactly-once 语义与第一节特性矩阵中的标注一致。在close()时连接器会关闭 Firestore 客户端并释放资源若关闭失败会抛出带错误码FIRESTORE-01Close Firestore client failed的FirestoreConnectorException错误码定义见 FirestoreConnectorErrorCode.java。五、字段类型映射SeaTunnel 类型 → Firestore 类型连接器将 SeaTunnel 类型转换为 Firestore 文档字段值完整映射关系如下表SeaTunnel 类型Firestore 值TINYINTintegerSMALLINTintegerINTintegerBIGINTintegerFLOATdoubleDOUBLEdoubleDECIMALdecimal valueSTRINGstringBOOLEANbooleanBYTESblobDATEdateUTC 当日零点TIMESTAMPtimestampARRAYarrayMAPmapNULLnull序列化的源码实现上述映射在 DefaultSeaTunnelRowSerializer.java 的convert方法中逐类型实现几个值得注意的细节TINYINT/SMALLINT/INT统一转为intValue()BIGINT转为longValue()FLOAT与DOUBLE都转为double写入DECIMAL以BigDecimal原样写入Firestore 支持 decimal valueBYTES通过Blob.fromBytes(...)转为 Firestore 的 BlobDATE以LocalDate.atStartOfDay(ZoneOffset.UTC)转换为 UTC 当日零点的DateTIMESTAMP以Timestamp.of(...)转为 Firestore TimestampARRAY递归转换每个元素为ListObjectMAP递归转换每个 Value 为对应的 Firestore 值字段值为null时直接映射为 null。序列化时上游 SeaTunnel Schema 的字段名会直接成为 Firestore 文档的字段名见serialize方法中data.put(seaTunnelRowType.getFieldName(index), ...)的逻辑。六、使用边界与注意事项Notes以下限制直接决定作业设计与数据模型规划请务必在开发前评估仅提供 Sink当前连接器只有写入端没有 GoogleFirestore Source 连接器单 Sink 单集合每个 Sink 块写入一个配置好的 Collection不会针对多表输入自动切换集合——多集合写入请为每个集合配置一个 Sink 块文档 ID 自动生成由于使用add追加写入文档 ID 由 Firestore 自动生成。如需确定性文档 ID请在进入本 Sink 之前使用其他连接器或 Transform 预处理不识别 CDC 语义连接器不会把UPDATE/DELETE行类型解释为 CDC 操作——每一行都会触发一次 Firestoreadd产生一个新文档凭据必须 Base64不要把原始 Service Account JSON 直接放在credentials中字段名继承上游 SeaTunnel Schema 的字段名会成为 Firestore 文档字段名checkpoint 语义连接器同时支持BATCH与STREAMING两种作业模式但由于write()逐行调用add且无内存缓冲checkpoint 完成不意味着之前所有行已成功写入 Firestore。七、任务示例Task Example7.1 批处理写入全类型字段验证以下配置使用FakeSource构造一条包含全部支持数据类型的行并通过 GoogleFirestore Sink 写入env { parallelism 1 job.mode BATCH } source { FakeSource { schema { fields { c_map mapstring, string c_array arrayint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_decimal decimal(30, 8) c_null null c_bytes bytes c_date date c_timestamp timestamp } } rows [ { kind INSERT fields [{a: b}, [10], c_string, true, 117, 15987, 56387395, 7084913402530365000, 1.23, 1.23, 2924137191386439303744.39292216, null, bWlJWmo, 2023-04-22, 2023-04-22T23:20:58] } ] } } sink { GoogleFirestore { project_id dummy-project collection dummy-collection credentials base64-service-account-json } }该配置与仓库 E2E 测试使用的 fake_to_google_firestore.conf 一致。对应的端到端测试 GoogleFirestoreIT.java 会执行作业后读取 Firestore断言写入文档字段数为 15并逐项校验15987Lsmallint、56387395Lint、2924137191386439303744.39292216decimal、Blob.fromBytes(...)bytes以及 Timestamp 等转换结果。该测试默认Disabled因为需要真实的 Google Firestore 数据库环境才能运行。7.2 流式写入配合 checkpoint 间隔以下配置演示STREAMING模式下、指定 checkpoint 间隔的写入env { parallelism 1 job.mode STREAMING checkpoint.interval 30000 } source { FakeSource { row.num 100 schema { fields { c_string string c_int int c_timestamp timestamp } } plugin_output firestore_stream } } sink { GoogleFirestore { plugin_input firestore_stream project_id my-gcp-project collection events credentials base64-service-account-json } }示例要点job.mode STREAMING且checkpoint.interval 30000单位毫秒开启 30 秒周期的 checkpoint通过plugin_output/plugin_input显式串联 FakeSource 与 GoogleFirestore Sink结合第四节所述的无缓冲逐行add模型请注意 checkpoint 不代表数据已达 Firestore。八、使用流程小结准备凭据在 Google Cloud 创建 Service Account 并下载 JSON执行base64 -w 0 service-account.jsonmacOS 用base64 service-account.json | tr -d \n得到配置值或确保运行环境可解析 Application Default Credentials如设置GOOGLE_APPLICATION_CREDENTIALS确认目标确定project_id与目标 Collection并按需在 Google Cloud 预先创建涉及查询的索引编写作业配置按第七节模板配置env/source/sink其中 Sink 块至少包含project_id与collection可选credentials运行与验证使用 SeaTunnel 引擎Zeta、Spark 或 Flink提交作业可在 Firestore 控制台或通过客户端查询确认文档已写入、类型映射符合第五节表格。九、版本演进参考依据 connector-google-firestore 变更日志2.3.2新增 GoogleFirestore Sink 连接器Feature2.3.4移除对SeaTunnelSink::getConsumedType的使用并标记废弃2.3.9允许将指标信息关联到逻辑计划节点2.3.10改进 Firestore 选项。如需进一步了解 Sink 公共选项请阅读 Sink Common Options若想基于本文示例扩展更多连接器用法可参考仓库中的 connector-google-firestore 模块 及其 E2E 测试。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

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

较早相关资讯

最新相关资讯

基于SpringBoot+SSM的校园防诈骗宣传平台设计与实现 2026/9/17 3:48:58

基于SpringBoot+SSM的校园防诈骗宣传平台设计与实现

2. 核心功能模块设计:用户端与管理员端2.1 用户端功能拆解:从“被动看”到“主动防”用户端是整个平台的流量入口,也是防骗教育真正落地的场景。在设计时,我把它拆成四大块:资讯浏览、案例学习、在线答题、留言反馈。每…

阅读更多 →
grub> 命令行救援指南:Linux 引导故障修复与预防 2026/9/17 3:48:58

grub> 命令行救援指南:Linux 引导故障修复与预防

开机之后没看到熟悉的桌面或者登录界面,屏幕上顶着一行grub>或者grub rescue>,下面还跟着一句minimal bash-like line editing is supported。第一次遇到的人十有八九会慌,以为系统挂了,其实大部分情况下数据都还在&#xf…

阅读更多 →
线性定常系统参数辨识:从Hankel矩阵到嵌入式部署的完整MATLAB实践 2026/9/17 3:48:58

线性定常系统参数辨识:从Hankel矩阵到嵌入式部署的完整MATLAB实践

简介:本资源是面向自动化、控制工程及系统建模方向高年级本科生与研究生的MATLAB实践程序包,聚焦线性定常系统的参数辨识核心问题,覆盖阶次已知与未知两类典型场景下的差分方程建模与参数估计方法。压缩包共9个文件,含4个功能完整…

阅读更多 →
Home Assistant 实战:input_button.reload 动作详解——不重启即可重载 YAML 中的按钮辅助组件 2026/9/17 3:48:58

Home Assistant 实战:input_button.reload 动作详解——不重启即可重载 YAML 中的按钮辅助组件

Home Assistant 实战:input_button.reload 动作详解——不重启即可重载 YAML 中的按钮辅助组件 【免费下载链接】home-assistant.io :blue_book: Home Assistant User documentation 项目地址: https://gitcode.com/GitHub_Trending/ho/home-assistant.io 本…

阅读更多 →
从“Helloword”到Hello World:编程第一课的环境验证与工具链价值 2026/9/17 3:48:58

从“Helloword”到Hello World:编程第一课的环境验证与工具链价值

写代码这么多年,有一个词,几乎每一位程序员都敲过,甚至闭着眼都能打出来——"Hello, World"。但有意思的是,在搜索引擎的热搜词里,它的拼写却是"Helloword",没有逗号,没有空…

阅读更多 →
STM32+RS485+MODBUS RTU精准时序通信实战 2026/9/17 3:45:58

STM32+RS485+MODBUS RTU精准时序通信实战

简介:本资源是一套基于STM32平台实现Modbus-RTU通信协议的完整双模工程,面向嵌入式初学者与中级开发者,解决工业现场RS485主从设备协同通信的典型开发难题。项目支持主机轮询与按键切换从机两种工作模式,具备地址01–03从机数据读…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

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

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