ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SeaTunnel MongoDB Sink Connector 完全指南:从数据写入到事务与幂等写入

SeaTunnel MongoDB Sink Connector 完全指南:从数据写入到事务与幂等写入 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本指南以 SeaTunnel 仓库中 docs/en/connector-v2/sink/MongoDB.md 为骨架结合 connector-mongodb 模块源码系统讲解 SeaTunnel MongoDB Sink 连接器的使用方式支持引擎与依赖、数据类型映射、全部 Sink 配置参数、Buffer Flush 触发机制、事务提交模式、基于 upsert 的幂等写入以及完整的可运行配置示例。读完本文你将能独立编写将任意 SeaTunnel 数据源写入 MongoDB 的同步任务并能根据业务需要正确选配transaction、upsert-enable、primary-key等关键参数。支持引擎与核心特性MongoDB Sink Connector 可运行于以下引擎SparkFlinkSeaTunnel Zeta其核心特性包括exactly-once精确一次语义通过 checkpoint 与事务/幂等写入配合实现详见 connector-v2-features。CDC 写入支持可承接上游 CDC 变更数据流INSERT / UPDATE / DELETE详见 connector-v2-features。提示如果希望使用 CDC 写入特性官方建议开启upsert-enable配置以便将 INSERT/UPDATE 统一转换为 upsert 语义避免重复数据处理时触发主键冲突。连接器描述MongoDB Connector 提供从 MongoDB 读取和写入数据的能力。本文档聚焦于如何配置 MongoDB 连接器以向 MongoDB 写入数据Sink 侧。从源码结构看connector-mongodb模块同时包含source与sink两大包sink/MongodbSink、MongodbWriter、MongodbWriterOptions、MongoKeyExtractor等source/MongodbSource、MongodbReader、MongodbSplitEnumerator等。Sink 侧的入口类是 MongodbSink.java它通过AutoService(SeaTunnelSink.class)注册为 SeaTunnel 官方 Sink 插件插件名称为MongoDB定义于 MongodbConfig.java 中的CONNECTOR_IDENTITY常量。依赖与数据源信息使用 MongoDB 连接器需要引入以下依赖可通过install-plugin.sh脚本或 Maven 中央仓库下载数据源支持版本依赖MongoDBuniversalseatunnel-connectors-v2/connector-mongodbMaven 坐标org.apache.seatunnel:seatunnel-connectors-v2:connector-mongodb在编译期该模块还依赖mongodb-driver-sync、bson等 MongoDB 官方驱动库见 connector-mongodb/pom.xml。数据类型映射Seatunnel 类型与 MongoDB BSON 类型Sink 写入时会经过 RowDataToBsonConverters.java 完成 Seatunnel 行数据到 BSON 文档的转换。官方映射关系如下Seatunnel 数据类型MongoDB BSON 类型STRINGObjectIdSTRINGStringBOOLEANBooleanBINARYBinaryINTEGERInt32TINYINTInt32SMALLINTInt32BIGINTInt64DOUBLEDoubleFLOATDoubleDECIMALDecimal128DateDateTimestampTimestamp[Date]ROWObjectARRAYArray使用注意Tips日期精度差异使用 SeaTunnel 的Date与Timestamp类型写入 MongoDB 时两者最终都会生成 MongoDB 的Date类型但精度不同——Date产生秒级精度Timestamp产生毫秒级精度。从源码看RowDataToBsonConverters.java 中DATE类型将LocalDate转成当天零点的BsonDateTime毫秒时间戳而TIMESTAMP类型将LocalDateTime直接转成毫秒时间戳。DECIMAL 精度上限使用DECIMAL类型时最大范围不能超过 34 位数字即应使用decimal(34, 18)。源码中DECIMAL通过BsonDecimal128(new Decimal128(...))转换RowDataToBsonConverters.java受 MongoDBDecimal128类型的 34 位精度限制。此外还有两个细节值得注意STRING 的特殊 JSON 解析当 STRING 值以{开头、以}结尾且包含_value字段时会尝试按 Extended JSON 解析出内部 BSON 值解析失败则回退为普通字符串存储RowDataToBsonConverters.java。MAP 的 key 必须为字符串MAP 类型转换要求 key 类型必须是 STRING否则抛出UNSUPPORTED_OPERATION异常RowDataToBsonConverters.java。Sink 参数详解所有 Sink 参数定义在 MongodbConfig.java 的 sink parameters 区域与文档表格一一对应参数名类型是否必填默认值说明uriString是-MongoDB 标准连接 URI例如mongodb://user:passwordhosts:27017/database?readPreferencesecondaryslaveOktruedatabaseString是-要读写数据的 MongoDB 数据库名collectionString是-要读写数据的 MongoDB 集合名schemaString是-MongoDB BSON 与 Seatunnel 数据结构映射buffer-flush.max-rowsString否1000每次批量请求缓冲的最大行数buffer-flush.intervalString否30000批量请求的最大缓冲间隔单位毫秒retry.maxString否3写库失败时的最大重试次数retry.intervalDuration否1000写库失败后的重试间隔单位毫秒upsert-enableBoolean否false是否以 upsert 模式写入文档primary-keyList否-upsert/update 的主键格式如[id,name,...]transactionBoolean否false是否在 MongoSink 中使用事务要求 MongoDB 4.2common-options否-Sink 插件通用参数详见 Sink Common Options从 MongodbSink.java 的prepare方法可以看出参数解析逻辑uri、database、collection三者缺一不可同时缺失时会静默不构建配置其余参数均为可选读取后通过 Builder 模式组装成 MongodbWriterOptions。Tips重要提示刷盘时机由三参数联合控制MongoDB Sink 的数据刷盘逻辑由buffer-flush.max-rows、buffer-flush.interval和checkpoint.interval三个参数共同控制任一条件满足即触发刷盘。从源码看MongodbWriter.java 的write()方法中非事务模式下每当isOverMaxBatchSizeLimit()缓冲行数达到max-rows或isOverMaxBatchIntervalLimit()距上次发送超过interval毫秒满足其一就调用doBulkWrite()同时prepareCommit()checkpoint 触发也会强制刷盘。upsert-key历史兼容兼容历史参数upsert-key。如果设置了upsert-key请不要再设置primary-key。源码中 MongodbConfig.java 将primary-key定义为带 fallback keyupsert-key的选项MongodbSink.java 会优先读取primary-key若未配置则回退读取upsert-key。快速上手创建 MongoDB 数据同步任务以下示例演示如何创建一个将随机生成数据写入 MongoDB 数据库的数据同步任务# Set the basic configuration of the task to be performed env { parallelism 1 job.mode BATCH checkpoint.interval 1000 } source { FakeSource { row.num 2 bigint.min 0 bigint.max 10000000 split.num 1 split.read-interval 300 schema { fields { c_bigint bigint } } } } sink { MongoDB{ uri mongodb://user:password127.0.0.1:27017 database test collection test schema { fields { _id string c_bigint bigint } } } }要点解析schema.fields中的字段名与类型必须与上游数据对应_id字段可显式声明用于指定文档主键env中设置checkpoint.interval 1000毫秒它与buffer-flush.*共同决定写入 MongoDB 的刷盘节奏FakeSource是内置模拟数据源可参考 connector-fake 模块了解其全部参数如row.num、split.num等。参数深度解读MongoDB 数据库连接 URI 示例未认证单节点连接mongodb://127.0.0.0:27017/mydb副本集连接mongodb://127.0.0.0:27017/mydb?replicaSetxxx带认证的副本集连接mongodb://admin:password127.0.0.0:27017/mydb?replicaSetxxxauthSourceadmin多节点副本集连接mongodb://127.0.0.1:27017,127.0.0.2:27017,127.0.0.3:27017/mydb?replicaSetxxx分片集群连接通过 mongos 接入URI 与普通连接一致mongodb://127.0.0.0:27017/mydb多 mongos 连接mongodb://192.168.0.1:27017,192.168.0.2:27017,192.168.0.3:27017/mydb注意URI 中的用户名和密码在拼接进连接字符串之前必须进行 URL 编码URL-encode否则特殊字符如、:、/会导致连接解析失败。Buffer Flush 刷盘机制示例配置sink { MongoDB { uri mongodb://user:password127.0.0.1:27017 database test_db collection users buffer-flush.max-rows 2000 buffer-flush.interval 1000 schema { fields { _id string id bigint status string } } } }底层实现原理MongoDB Sink 的写入采用内存批量缓冲 驱动 BulkWrite的方式。核心逻辑在 MongodbWriter.java 的doBulkWrite()待写文档先缓存在bulkRequestsListWriteModelBsonDocument中不逐条写入满足刷盘条件后调用 MongoDB 驱动的bulkWrite(bulkRequests, new BulkWriteOptions().ordered(true))以有序批量方式一次性提交批量写入失败时按retry.max重试重试间隔按retry.interval * (i 1)递增退避i 为重试次数从 0 开始重试耗尽后抛出MongodbConnectorException(WRITER_OPERATION_FAILED)任务进入失败处理流程close()时也会执行最后一次刷盘确保残留缓冲数据落库。因此调大buffer-flush.max-rows、buffer-flush.interval可显著减少网络往返、提升吞吐但会增大单批内存占用与故障恢复时的重复数据范围调小则降低延迟、提升实时性适合 CDC / 流式场景。为什么不建议盲目使用事务尽管 MongoDB 自 4.2 版本起已完整支持多文档事务但这并不意味着应该随意滥用。事务本质上是锁 节点协调 额外开销 性能损耗的叠加。使用事务的原则应当是能不用的地方尽量不用。通过合理设计系统例如保证写入顺序、使用幂等键、接受最终一致可以很大程度上避免对事务的依赖。事务模式的源码实现当transaction true时MongodbSink.java 会启用SinkAggregatedCommitter与状态序列化器。其提交过程在 MongodbSinkAggregatedCommitter.java 中实现先client.startSession()开启客户端会话将缓冲文档封装为DocumentBulk使用clientSession.withTransaction(...)在事务内提交事务选项为readPreferenceprimary、readConcernLOCAL、writeConcernMAJORITY每个 DocumentBulk 默认缓冲上限为 1024 个文档BUFFER_SIZE 1024但事务模式下不遵守该上限因为同一批数据必须在一个事务中整体提交提交失败的数据会返回给引擎由 checkpoint 机制在后续恢复时重试。另外事务模式还要求 MongoDB 部署支持事务4.2 副本集或分片集群因此只有业务确有跨文档原子性需求例如 CDC 批量变更必须原子落库时才应开启。幂等写入Idempotent Writes通过指定明确的主键并使用 upsert 方法可以实现exactly-once 写入语义。当配置中定义了primary-key且开启upsert-enable时MongoDB Sink 将使用upsert 语义而非普通的 INSERT 语句把primary-key声明的键组合作为 MongoDB 文档的过滤条件以 upsert 模式写入从而保证幂等。为什么重要任务失败时SeaTunnel 会从最近一次成功的 checkpoint 恢复并重新处理恢复过程中可能产生重复消息。强烈建议使用 upsert 模式它有助于避免记录被重新处理时违反数据库主键约束、产生重复数据。源码验证RowDataDocumentSerializer.java 中按RowKind构造写入模型INSERT开启 upsert 时走UpdateOneModel upsert(true)$set更新、不存在则插入否则走InsertOneModelUPDATE_AFTER开启 upsert 时同样走 upsert 更新否则走普通UpdateOneModelDELETE走DeleteOneModel按主键过滤条件删除。其中 upsert 的过滤条件由 MongoKeyExtractor.java 从每条 BSON 文档中抽取primary-key指定字段生成并通过Filters.and(Filters.eq(...))组合成等值过滤条件。示例配置sink { MongoDB { uri mongodb://user:password127.0.0.1:27017 database test_db collection users upsert-enable true primary-key [name,status] schema { fields { _id string name string status string } } } }实践建议在 CDC 同步、实时数仓增量写入等场景中upsert-enable true配合合理的primary-key是兼顾正确性与性能的推荐组合transaction true仅在确有跨文档原子性要求的场景下开启。Changelog2.2.0-beta新增 MongoDB Source Connector。2.3.1-release重构 MongoDB Source Connector。Next VersionMongoDB 支持 CDC Sink关联 PR 见仓库 release-note 与 issue 记录。进一步探索连接器完整源码seatunnel-connectors-v2/connector-mongodb连接器工厂与插件注册MongodbSinkFactory.java连接器特性总览connector-v2-features.md作业配置通用语法config.md端到端测试用例connector-mongodb-e2e赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel MongoDB Sink Connector 完全指南从批量写入到事务与幂等 Upsert 的实战配置SeaTunnel MongoDB Sink Connector 完全指南从批量写入到事务与幂等 Upsert 的实战配置 导读 本文是 SeaTunnel数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Neo4j Sink Connector 使用指南从逐条写入到 UNWIND 批量写入SeaTunnel Neo4j Sink Connector 使用指南从逐条写入到 UNWIND 批量写入 SeaTunnel 的 Neo4j Sink 插件数据工程大数据批处理流处理SeaTunnel Lance Sink Connector 完整指南将海量数据写入 Lance 数据集SeaTunnel Lance Sink Connector 完整指南将海量数据写入 Lance 数据集 本篇指南围绕 SeaTunnel 官方文档 docs数据集成ETL大数据批处理流处理变更数据捕获上一篇Ollama本地大模型革命5分钟部署Llama 3.2全攻略下一篇Cycle.js 版本演进全记录从 Cycle Nested 到 Synthetic DOM Events 的迁移实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表