语义落地)
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-enableprimary-key将精确一次语义落到 MongoDB。概述MongoDB Sink 连接器在 SeaTunnel 中的定位MongoDB Sink 连接器负责将 SeaTunnel 的行记录SeaTunnelRow写入 MongoDB 集合每行记录都会被转换为 BSON 文档发送到配置的database与collection。连接器的插件标识为MongoDB见 MongodbBaseOptions.java 中的CONNECTOR_IDENTITY整体 Sink 实现见 MongodbSink.java。支持的引擎SparkFlinkSeaTunnel Zeta关键特性exactly-once 精准一次写入定时刷新仅 Zeta 引擎CDC变更数据捕获支持多表写入使用提示如果希望使用 CDC 写入功能建议启用upsert-enable配置项。从源码看RowDataDocumentSerializer.java 会依据RowKind将 CDC 的 INSERT、UPDATE_AFTER、DELETE 分别转换为对应的WriteModel而 MongodbWriter.java 会过滤掉UPDATE_BEFORE记录避免冗余写入。启用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.javaINSERT开启 upsert 时使用UpdateOneModel带UpdateOptions().upsert(true)否则使用InsertOneModelUPDATE_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-mongodbMaven 中央仓库 /install-plugin.sh安装到本地后请确认config/plugin_config中已注册connector-mongodb并在connector-jar目录下存在对应的 JAR即可在作业配置中以MongoDB { ... }块引用。数据类型映射下表展示了 SeaTunnel 数据类型到 MongoDB BSON 类型的映射关系。SeaTunnel 数据类型MongoDB BSON 类型STRINGObjectIdSTRINGStringBOOLEANBooleanBINARYBinaryINTEGERInt32TINYINTInt32SMALLINTInt32BIGINTInt64DOUBLEDoubleFLOATDoubleDECIMALDecimal128DateDateTimestampTimestamp / DateROWObjectARRAYArray映射的底层实现在 RowDataToBsonConverters.java可从源码确认以下细节整型家族TINYINT、SMALLINT、INT统一转换为BsonInt32BIGINT转换为BsonInt64见 RowDataToBsonConverters.java。浮点家族FLOAT、DOUBLE均转换为BsonDouble见 RowDataToBsonConverters.java。二进制BYTESBINARY转换为BsonBinary。日期时间DATE按系统时区的当日零点转成 epoch 毫秒TIMESTAMP则直接按本地时间转换。两者最终都落为 BSONDate。小数DECIMAL使用Decimal128承载转换时会带入字段定义的 precision 与 scale见 RowDataToBsonConverters.java。复合类型ARRAY递归转换为BsonArrayROW转换为内嵌BsonDocumentObject。提示使用 SeaTunnel 将Date和Timestamp类型写入 MongoDB 时结果都是 BSONDate类型但精度不同SeaTunnel 的Date为秒级日期粒度精度Timestamp为毫秒级精度。在 SeaTunnel 中使用DECIMAL类型时最大精度不能超过 34 位对应Decimal128的能力上限。建议使用decimal(34, 18)以满足支持的精度与标度。Sink 参数说明参数名称类型是否必填默认值描述uriString是-MongoDB 标准连接 URI例如mongodb://user:passwordhosts:27017/database?readPreferencesecondaryslaveOktrue。更多示例请参考下文「连接 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-rowsInt否1000每次批量写入请求的最大缓存行数。buffer-flush.intervalLong否30000批量写入请求的最大时间间隔毫秒。retry.maxInt否3写入失败时的最大重试次数。retry.intervalLong否1000写入失败后的重试间隔毫秒。upsert-enableBoolean否false是否启用 upsert 模式写入。开启时需要同时配置primary-key。primary-keyList否-用于 upsert 或更新的主键格式为[id,name,...]。transactionBoolean否false是否在 MongoSink 中启用事务需要 MongoDB 4.2。data_save_modeEnum否APPEND_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默认false、transaction默认false、data_save_mode默认APPEND_DATA。提示MongoDB Sink 的连接器级数据刷新由三个参数共同控制buffer-flush.max-rows、buffer-flush.interval和checkpoint.interval。任一条件触发都会立刻刷写。从源码看MongodbWriter.write()在非事务模式下会检查isOverMaxBatchSizeLimit()缓存行数达到bulkActions与isOverMaxBatchIntervalLimit()距上次发送超过batchIntervalMs满足其一即触发doBulkWrite()见 MongodbWriter.java。兼容历史参数upsert-key作为primary-key的回退名。若已设置upsert-key请勿同时设置primary-key。源码中通过Options.key(primary-key).withFallbackKeys(upsert-key)实现见 MongodbSinkOptions.java。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:password127.0.0.1:27017 database test collection test } }将上述配置保存为作业配置文件后使用bin/seatunnel.sh --config 配置文件路径即可提交运行Zeta 引擎下默认job.mode支持BATCH与STREAMING两种模式。多表写入当上游记录携带表元数据例如 CDC 场景或FakeSource的tables_configs时database和collection可以使用占位符。常用占位符包括${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?retryWritestrue 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?replicaSetxxx带认证的副本集连接mongodb://admin:password127.0.0.1:27017/mydb?replicaSetxxxauthSourceadmin多节点副本集连接mongodb://192.168.0.1:27017,192.168.0.2:27017,192.168.0.3:27017/mydb?replicaSetxxx分片集群连接通过一个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:password127.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 驱动bulkWriteBulkWriteOptions().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.LOCAL、WriteConcern.MAJORITY见 MongodbSinkAggregatedCommitter.java。这也是上文所述「事务模式通过 checkpoint 提交、因此禁用定时刷新」的原因——事务边界必须与 checkpoint 边界保持一致。幂等写入Idempotent Writes把 exactly-once 落到 MongoDB通过定义明确的主键并启用upsert模式可以实现精准一次写入exactly-once语义。当配置中定义了primary-key且启用了upsert-enableMongoDB Sink 将使用 Upsert 语义而非普通 INSERT 语句。SeaTunnel 会将定义的主键作为 MongoDB 的复合主键在 Upsert 模式下写入以确保幂等性。若作业在运行过程中失败SeaTunnel 会从上一个成功的 checkpoint 恢复并重新处理数据这可能导致重复数据。强烈建议启用 Upsert 模式以避免主键冲突或重复插入。sink { MongoDB { uri mongodb://user:password127.0.0.1:27017 database test_db collection users upsert-enable true primary-key [name, status] } }在该配置下每次写入都会以name、status两列作为过滤条件执行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),仅供参考