SeaTunnel MongoDB Sink 连接器实战指南:从批量写入到精确一次(Exactly-Once)语义落地
2026/9/19 0:55:16 网站建设 项目流程

SeaTunnel MongoDB Sink 连接器实战指南:从批量写入到精确一次(Exactly-Once)语义落地

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

本文以 SeaTunnel 仓库中 MongoDB Sink 连接器文档 为主体,结合seatunnel-connectors-v2/connector-mongodb模块源码,系统讲解如何将 SeaTunnel 行记录写入 MongoDB 集合。读完本文,你将掌握 MongoDB Sink 的全部配置项与默认值、追加/Upsert 两种写入语义的区别、数据类型到 BSON 的映射规则、Zeta 引擎定时刷新、多表写入占位符,以及如何借助upsert-enable+primary-key将精确一次语义落到 MongoDB。

概述:MongoDB Sink 连接器在 SeaTunnel 中的定位

MongoDB Sink 连接器负责将 SeaTunnel 的行记录(SeaTunnelRow)写入 MongoDB 集合:每行记录都会被转换为 BSON 文档,发送到配置的databasecollection。连接器的插件标识为MongoDB(见 MongodbBaseOptions.java 中的CONNECTOR_IDENTITY),整体 Sink 实现见 MongodbSink.java。

支持的引擎

  • Spark
  • Flink
  • SeaTunnel Zeta

关键特性

  • exactly-once 精准一次写入
  • 定时刷新(仅 Zeta 引擎)
  • CDC(变更数据捕获)
  • 支持多表写入

使用提示

  1. 如果希望使用 CDC 写入功能,建议启用upsert-enable配置项。从源码看,RowDataDocumentSerializer.java 会依据RowKind将 CDC 的 INSERT、UPDATE_AFTER、DELETE 分别转换为对应的WriteModel,而 MongodbWriter.java 会过滤掉UPDATE_BEFORE记录,避免冗余写入。
  2. 启用transaction与 Zeta 定时刷新互斥,请在同一个作业中选择其中一种模式;混用会导致定时刷新被静默关闭。这一点在源码中有直接体现:MongodbWriter构造函数仅在!transaction时才会context.registerFlushAction(this::timerFlush)注册定时刷新动作(见 MongodbWriter.java)。

两种写入语义:追加写入与 Upsert 写入

连接器支持两种写入语义:

  • 追加写入(Append):每行生成一条新文档,通过InsertOneModel写入。性能最好,但在失败重试场景下不幂等——同一行数据可能被插入多条。
  • Upsert 写入:当upsert-enable = true且配置了primary-key时,连接器会把主键字段作为 MongoDB 的_id(或复合_id)并以 upsert 方式写入。结合 checkpoint 恢复机制,可以实现 at-least-once + 幂等重试,这是把 exactly-once 落到 MongoDB 的标准做法。

源码层面的映射逻辑在 RowDataDocumentSerializer.java:

  • INSERT:开启 upsert 时使用UpdateOneModel(带UpdateOptions().upsert(true)),否则使用InsertOneModel
  • UPDATE_AFTER:开启 upsert 时同样走 upsert 模型,否则走普通UpdateOneModel(仅更新、不插入);
  • DELETE:使用DeleteOneModel

upsert 的过滤条件由 MongoKeyExtractor.java 根据primary-key从 BSON 文档中提取主键字段生成,最终在generateFilter中组合为Filters.and(Filters.eq(...))形式(见 RowDataDocumentSerializer.java)。

缓存、重试以及可选的事务都可以通过下文的配置项进行调节。

依赖与安装

要使用 MongoDB 连接器,需要以下依赖。可以通过install-plugin.sh下载,也可以从 Maven 中央仓库获取(Artifact 坐标为org.apache.seatunnel:connector-mongodb)。

数据源支持版本依赖
MongoDB通用版本connector-mongodb(Maven 中央仓库 /install-plugin.sh

安装到本地后,请确认config/plugin_config中已注册connector-mongodb,并在connector-jar目录下存在对应的 JAR,即可在作业配置中以MongoDB { ... }块引用。

数据类型映射

下表展示了 SeaTunnel 数据类型到 MongoDB BSON 类型的映射关系。

SeaTunnel 数据类型MongoDB BSON 类型
STRINGObjectId
STRINGString
BOOLEANBoolean
BINARYBinary
INTEGERInt32
TINYINTInt32
SMALLINTInt32
BIGINTInt64
DOUBLEDouble
FLOATDouble
DECIMALDecimal128
DateDate
TimestampTimestamp / Date
ROWObject
ARRAYArray

映射的底层实现在 RowDataToBsonConverters.java,可从源码确认以下细节:

  • 整型家族TINYINTSMALLINTINT统一转换为BsonInt32BIGINT转换为BsonInt64(见 RowDataToBsonConverters.java)。
  • 浮点家族FLOATDOUBLE均转换为BsonDouble(见 RowDataToBsonConverters.java)。
  • 二进制BYTES(BINARY)转换为BsonBinary
  • 日期时间DATE按系统时区的当日零点转成 epoch 毫秒;TIMESTAMP则直接按本地时间转换。两者最终都落为 BSONDate
  • 小数DECIMAL使用Decimal128承载,转换时会带入字段定义的 precision 与 scale(见 RowDataToBsonConverters.java)。
  • 复合类型ARRAY递归转换为BsonArrayROW转换为内嵌BsonDocument(Object)。

提示

  1. 使用 SeaTunnel 将DateTimestamp类型写入 MongoDB 时,结果都是 BSONDate类型,但精度不同:SeaTunnel 的Date为秒级(日期粒度)精度,Timestamp为毫秒级精度。
  2. 在 SeaTunnel 中使用DECIMAL类型时,最大精度不能超过 34 位(对应Decimal128的能力上限)。建议使用decimal(34, 18)以满足支持的精度与标度。

Sink 参数说明

参数名称类型是否必填默认值描述
uriString-MongoDB 标准连接 URI,例如mongodb://user:password@hosts:27017/database?readPreference=secondary&slaveOk=true。更多示例请参考下文「连接 URI 详解」。
databaseString-要写入的 MongoDB 数据库名称。配置多表同步时,可使用占位符${database_name},例如:database = "${database_name}_test_database"
collectionString-要写入的 MongoDB 集合名称。配置多表同步时,可使用${database_name}${schema_name}${table_name}等占位符,例如:collection = "${database_name}_${schema_name}_${table_name}_check"
buffer-flush.max-rowsInt1000每次批量写入请求的最大缓存行数。
buffer-flush.intervalLong30000批量写入请求的最大时间间隔(毫秒)。
retry.maxInt3写入失败时的最大重试次数。
retry.intervalLong1000写入失败后的重试间隔(毫秒)。
upsert-enableBooleanfalse是否启用 upsert 模式写入。开启时需要同时配置primary-key
primary-keyList-用于 upsert 或更新的主键,格式为["id","name",...]
transactionBooleanfalse是否在 MongoSink 中启用事务(需要 MongoDB 4.2+)。
data_save_modeEnumAPPEND_DATAMongoDB 集合的数据写入模式:DROP_DATA表示写入前清空集合;APPEND_DATA表示追加写入;ERROR_WHEN_DATA_EXISTS表示集合已有数据时直接报错。
common-options--通用 Sink 插件参数,详见 Sink Common Options。

以上参数的定义与默认值可在 MongodbSinkOptions.java 中逐一核对,例如buffer-flush.max-rows默认 1000、buffer-flush.interval默认 30000ms、retry.max默认 3、retry.interval默认 1000ms、upsert-enable默认falsetransaction默认falsedata_save_mode默认APPEND_DATA

提示

  1. MongoDB Sink 的连接器级数据刷新由三个参数共同控制:buffer-flush.max-rowsbuffer-flush.intervalcheckpoint.interval。任一条件触发都会立刻刷写。从源码看,MongodbWriter.write()在非事务模式下会检查isOverMaxBatchSizeLimit()(缓存行数达到bulkActions)与isOverMaxBatchIntervalLimit()(距上次发送超过batchIntervalMs),满足其一即触发doBulkWrite()(见 MongodbWriter.java)。
  2. 兼容历史参数upsert-key作为primary-key的回退名。若已设置upsert-key,请勿同时设置primary-key。源码中通过Options.key("primary-key").withFallbackKeys("upsert-key")实现(见 MongodbSinkOptions.java)。
  3. transaction选项与 Zeta 定时刷新互斥,请二选一。

Zeta 定时刷新(sink.flush.interval)

该引擎级能力仅由 Zeta 支持,Spark 和 Flink 不会注入FlushSignal记录。在 Zeta 中可以在env块配置sink.flush.interval,使未达到buffer-flush.max-rows的待处理 bulk 请求也能定时刷写出去。和buffer-flush.interval不同,引擎定时器不依赖新记录到达即可触发检查——即使数据流暂时停顿,待处理的 bulk 也会被定时刷出。

定时刷新仅在transaction = false时启用。MongoDB 事务模式通过 checkpoint 提交,因此会禁用定时刷新以保持事务边界。初始定时刷新实现提供至少一次(at-least-once)语义,不提供基于 2PC 的精确一次语义。启用 upsert 并使用确定性主键可使重试具备幂等性。

env { job.mode = "STREAMING" checkpoint.interval = 300000 sink.flush.interval = 5000 } sink { MongoDB { uri = "mongodb://127.0.0.1:27017" database = "test_db" collection = "users" buffer-flush.max-rows = 10000 transaction = false } }

该配置中,即使 10000 行阈值长期达不到,Zeta 引擎也会每 5000ms 发送一次刷新信号,MongodbWriter.timerFlush()收到信号后立即执行doBulkWrite()(见 MongodbWriter.java)。

快速上手:创建 MongoDB 数据同步任务

下面示例展示了一个将随机生成的数据写入 MongoDB 的任务:

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:password@127.0.0.1:27017" database = "test" collection = "test" } }

将上述配置保存为作业配置文件后,使用bin/seatunnel.sh --config <配置文件路径>即可提交运行(Zeta 引擎下默认job.mode支持BATCHSTREAMING两种模式)。

多表写入

当上游记录携带表元数据(例如 CDC 场景或FakeSourcetables_configs)时,databasecollection可以使用占位符。常用占位符包括${database_name}${schema_name}${table_name},例如collection = "${database_name}_${schema_name}_${table_name}_check"

source { FakeSource { tables_configs = [ { schema = { table = "testDatabase1.testSchema1.testTable1" fields { id = int value = string } } rows = [ { kind = INSERT fields = [1, "NEW"] } ] }, { schema = { table = "testDatabase2.testSchema2.testTable2" fields { id = int amount = "decimal(16, 1)" } } rows = [ { kind = INSERT fields = [1, 6.3] } ] } ] } } sink { MongoDB { uri = "mongodb://127.0.0.1:27017/test_db?retryWrites=true" database = "test_db" collection = "${database_name}_${schema_name}_${table_name}_check" } }

多表写入的能力由 MongodbSink.java 中实现的SupportMultiTableSink接口支撑,同时MongodbWriter实现了SupportMultiTableSinkWriter(见 MongodbWriter.java),可在同一作业中按表元数据动态解析目标集合。

参数详解:MongoDB 连接 URI 示例

无认证的单节点连接

mongodb://127.0.0.1:27017/mydb

副本集连接

mongodb://127.0.0.1:27017/mydb?replicaSet=xxx

带认证的副本集连接

mongodb://admin:password@127.0.0.1:27017/mydb?replicaSet=xxx&authSource=admin

多节点副本集连接

mongodb://192.168.0.1:27017,192.168.0.2:27017,192.168.0.3:27017/mydb?replicaSet=xxx

分片集群连接(通过一个mongos路由)

mongodb://mongos1.example.com:27017,mongos2.example.com:27017,mongos3.example.com:27017/mydb

多个 mongos 节点连接

mongodb://192.168.0.1:27017,192.168.0.2:27017,192.168.0.3:27017/mydb

注意:URI 中的用户名与密码在拼接前必须进行 URL 编码。

Buffer Flush 示例

sink { MongoDB { uri = "mongodb://user:password@127.0.0.1:27017" database = "test_db" collection = "users" buffer-flush.max-rows = 2000 buffer-flush.interval = 1000 } }

该配置表示:缓存满 2000 行,或距上次批量写入超过 1000ms,二者任一满足即触发一次批量写入。

批量写入与重试机制的底层实现

MongodbWriter.doBulkWrite()是批量刷写的核心(见 MongodbWriter.java):

  • 使用 MongoDB Java 驱动bulkWrite+BulkWriteOptions().ordered(true)按序写入;
  • 写入失败时按retry.max上限重试,第i次失败后休眠retryIntervalMs * (i + 1)毫秒,即重试间隔随次数线性退避;
  • 达到最大重试次数仍失败时抛出MongodbConnectorException(错误码WRITER_OPERATION_FAILED),作业据此进入失败处理与 checkpoint 恢复流程。

事务模式:为什么不推荐频繁使用事务?

虽然 MongoDB 自 4.2 版本起已完全支持多文档事务,但这并不意味着所有场景都应使用。事务意味着加锁、节点协调、额外往返和性能损耗。设计系统时应遵循的原则是:能不用事务就不要用事务。合理的系统设计可以在大多数情况下避免对事务的依赖。

如果确实启用transaction = true,写入路径会切换到事务模式:MongodbWriter.prepareCommit()不再直接刷写,而是把待写文档封装进DocumentBulk提交信息;随后由 MongodbSinkAggregatedCommitter.java 在 checkpoint 提交阶段使用clientSession.withTransaction(...)提交,事务选项为ReadPreference.primary()ReadConcern.LOCALWriteConcern.MAJORITY(见 MongodbSinkAggregatedCommitter.java)。这也是上文所述「事务模式通过 checkpoint 提交、因此禁用定时刷新」的原因——事务边界必须与 checkpoint 边界保持一致。

幂等写入(Idempotent Writes):把 exactly-once 落到 MongoDB

通过定义明确的主键并启用upsert模式,可以实现精准一次写入(exactly-once)语义。

当配置中定义了primary-key且启用了upsert-enable,MongoDB Sink 将使用 Upsert 语义而非普通 INSERT 语句。SeaTunnel 会将定义的主键作为 MongoDB 的复合主键,在 Upsert 模式下写入,以确保幂等性。

若作业在运行过程中失败,SeaTunnel 会从上一个成功的 checkpoint 恢复并重新处理数据,这可能导致重复数据。强烈建议启用 Upsert 模式,以避免主键冲突或重复插入。

sink { MongoDB { uri = "mongodb://user:password@127.0.0.1:27017" database = "test_db" collection = "users" upsert-enable = true primary-key = ["name", "status"] } }

在该配置下,每次写入都会以namestatus两列作为过滤条件执行UpdateOneModel(..., upsert(true)):文档存在则更新、不存在则插入,配合 checkpoint 恢复重放,同一行数据无论被处理多少次,最终集合中只有一份结果,从而在 at-least-once 的重试机制之上实现幂等,这是连接器精确一次语义的标准落地路径。

更新日志

MongoDB 连接器的历史变更记录(版本演进、功能新增与缺陷修复)见 connector-mongodb 更新日志,例如 2.3.3 版本引入的事务写入与 CDC Sink 支持、2.3.4 版本的 schema 主键/约束键配置支持等,可供升级排障时对照参考。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询