多亏了MongoDB的灵活性,我们早期业务跑得飞快,但后来数据量大了,麻烦也跟着来了。我们想把历史数据完整地存下来,还要能随时查,同时得跟上线上的变更。试过直接把MongoDB的变更流接到下游,可一旦遇到需要回刷、对比、或者做复杂的增量计算,就发现光靠MongoDB本身很难搞定。后来我们盯上了Hudi,把它当作全量历史数据存储层,从MongoDB迁过来,经历了不少磕磕绊绊。这篇文章就把我们走过的路、踩过的坑,还有最终能跑通的方案,原原本本讲给你听。
一、为什么要把MongoDB的数据搬到Hudi上
大概每个用MongoDB做核心库的团队,都会遇到“保存历史”的诉求。MongoDB本身不像关系型数据库那样容易做轻量级的快照恢复,虽然支持副本集和oplog,但oplog的保存时间很有限。我们的业务数据经常要回溯,比如用户改了手机号,我们要知道他以前用的什么号;订单状态变了好几次,我们想看看每次变化的时间点。这些都得靠完整的变更记录。
如果只把变更流存进Kafka,时间久了Kafka里的消息会过期,而且消息格式来回变,下游解析也麻烦。我们得有一个地方,既能放全量快照,又能把增量变更按主键合并进来,查的时候还得快。Hudi正好就是干这个的:它把数据组织成表,支持upsert(更新+插入),能记录每行数据的变化历史,还支持基于文件级别的索引,查询性能也不错。对我们来说,它就像一块“无限大的橡皮泥”,MongoDB里的数据怎么变,我们都能在Hudi里揉出想要的样子。
1.1 MongoDB同步到传统数仓的尴尬
以前我们用过定时拉全量数据到数仓的笨办法。每天凌晨跑一次MongoDB导出,把数据搬到Hive表里。这样能勉强回答“昨天是什么样”,但到了白天,业务暴增,导出的任务还没跑完,线上库就快被压垮了。更难受的是,一旦中间断了,整个表的数据就“对不上账”。
我们也试过直接用Debezium连接MongoDB的变更流,把变更事件打到Kafka,再消费进数仓。但Debezium给的是一条条change event,需要在数仓里自己搞一张“大宽表”,然后不断应用更新。做一次两次还行,数据量一大,主键更新冲突、重复消息、乱序到达就全都冒出来了。后来我们意识到,缺一个能统一处理“全量+增量”的存储层,这才把目光放在了Hudi上。
1.2 Hudi凭啥能接这个活
Hudi的核心能力是“对文件做增量管理”。它不像Hive那样每次重写整个分区,而是把更新写进小的文件组,再通过compaction合并成大文件。这样我们既能全量读,也能增量读,还能做时间旅行回到某个时刻。
对我们最有吸引力的,是Hudi的Copy-on-Write(COW)和Merge-on-Read(MOR)两种表类型。COW写的时候直接合并,读取简单,适合读多写少;MOR写的时候先记log文件,读的时候再合log,适合写多读少。我们的场景是写不少、读也算频繁,所以选用了MOR表。借助Spark的批量写入,我们可以在夜间把MongoDB全量数据导入Hudi,白天再用结构化流持续消费变更。
二、整体迁移思路:先复制,再持续追平
我们设计的迁移流程一共有四步。
第一步,打底。跑一个Spark作业,把MongoDB当前的全量数据按主键同步到Hudi表。第二步,追增量。与此同时,我们启动一个常驻的流式任务,从MongoDB的变更流里读取insert/update/delete操作,也写到同一张Hudi表。第三步,校验。写完之后,对比两边的主键集合和最新值,把差异补上。第四步,切换。线上应用不再直接访问MongoDB来查历史数据,而是改查Hudi(或者查Hudi同步到Doris的数据)。
这里最关键的一点是“全量”和“增量”不能打架。我们给Hudi表里设计了一个字段叫source_ts,表示这条记录最后一次变更的MongoDB时间戳。全量导入时,source_ts就是数据导出的那一刻;增量数据则用每条变更事件自身的clusterTime。写入Hudi时,我们用source_ts的大小来定先后,老时间戳不能覆盖新时间戳。这样就避免了并发写导致的“旧数据覆盖新数据”。
三、具体实现:从Spark到Hudi的完整示例
下面所有示例都用Java和Spark Structured Streaming,这个技术栈跟我们的生产环境一致。如果你用Python,思路也一样,照着改写就行。
3.1 环境准备
我们需要以下依赖:
<!-- pom.xml 里添加关键依赖 -->
<dependency>
<groupId>org.apache.hudi</groupId>
<artifactId>hudi-spark3.3-bundle_2.12</artifactId>
<version>0.14.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>3.3.2</version>
</dependency>
<dependency>
<groupId>org.mongodb.spark</groupId>
<artifactId>mongo-spark-connector_2.12</artifactId>
<version>10.2.0</version>
</dependency>
3.2 全量数据导入
我们用一个批量Spark作业读取MongoDB集合,然后写成Hudi表。这里用的是spark.read,从MongoDB读出的数据是DataFrame,里面每一行对应MongoDB里的一条文档。我们把它转成Hudi需要的格式,再调用df.write().format("hudi")写出去。
import org.apache.spark.sql.*;
import org.apache.spark.sql.types.*;
import org.apache.hudi.DataSourceWriteOptions;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.common.model.HoodieTableType;
public class FullSyncMongoToHudi {
public static void main(String[] args) {
SparkSession spark = SparkSession.builder()
.appName("mongo-full-sync-to-hudi")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog")
.getOrCreate();
// 1. 从MongoDB读取全量数据
Dataset<Row> mongoDf = spark.read()
.format("mongodb")
.option("connection.uri", "mongodb://user:pass@mongo-host:27017")
.option("database", "business_db")
.option("collection", "users")
.load();
// 2. 把MongoDB的_id转成字符串主键,并加一个source_ts字段
Dataset<Row> hudiDf = mongoDf
.withColumn("id", functions.col("_id").cast(DataTypes.StringType))
.withColumn("source_ts", functions.lit(System.currentTimeMillis()))
.drop("_id");
// 3. 写出Hudi表(MOR类型,主键为id,预合并字段为source_ts)
hudiDf.write()
.format("hudi")
.option(DataSourceWriteOptions.TABLE_TYPE_OPT_KEY(), HoodieTableType.MERGE_ON_READ.name())
.option(DataSourceWriteOptions.RECORDKEY_FIELD_OPT_KEY(), "id") // 主键字段
.option(DataSourceWriteOptions.PRECOMBINE_FIELD_OPT_KEY(), "source_ts") // 合并时根据哪个字段判断新旧
.option(DataSourceWriteOptions.PARTITIONPATH_FIELD_OPT_KEY(), "source_ts") // 这个示例按照时间做分区
.option(DataSourceWriteOptions.OPERATION_OPT_KEY(), DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL()) // 写入模式:更新+插入
.option(HoodieWriteConfig.TABLE_NAME, "hudi_users")
.mode(SaveMode.Append)
.save("/warehouse/hudi/users");
}
}
注意,我们这里用source_ts做分区字段,实际生产建议按日期yyyyMMdd分区,否则数据全堆在一个分区里,后续查询和clustering都会受到影响。上面代码里为了简化,直接用了毫秒时间戳,你上线前记得改成按天分区。
3.3 增量同步:用结构化流消费MongoDB变更流
在全量导入跑完后,我们要启动一个流式作业,监听MongoDB的变更流。MongoDB的change stream功能跟Spark的readStream配合得不错。我们每个批次读一批变更事件,然后把它们upsert到Hudi表。
import org.apache.spark.sql.*;
import org.apache.hudi.DataSourceWriteOptions;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.common.model.HoodieTableType;
import static org.apache.spark.sql.functions.*;
public class IncrementalSyncMongoToHudi {
public static void main(String[] args) throws Exception {
SparkSession spark = SparkSession.builder()
.appName("mongo-cdc-sync-to-hudi")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog")
.getOrCreate();
// 1. 从MongoDB的变更流中读取数据
Dataset<Row> changeStreamDf = spark.readStream()
.format("mongodb")
.option("connection.uri", "mongodb://user:pass@mongo-host:27017")
.option("database", "business_db")
.option("collection", "users")
.option("change.stream.publish.full.document.only", "true") // 只取更新后的完整文档
.option("change.stream.full.document", "updateLookup") // 更新时带上新文档内容
.load();
// 2. 将变更流解析成“Hudi行”的格式
// 变更流的格式为:{documentKey: {_id: ...}, fullDocument: {...}, operationType: "insert"|"update"|"delete", clusterTime: ...}
Dataset<Row> cdcDf = changeStreamDf
.select(
// 从documentKey里提取id字段
col("documentKey._id").cast("string").as("id"),
// 取出完整文档,如果没有则置空
col("fullDocument").as("doc"),
// 操作类型
col("operationType").as("op"),
// 变更时间(秒级,注意转换)
col("clusterTime").cast("long").as("source_ts")
);
// 3. 把嵌套文档转成键值对,方便Hudi存储
// 实际业务中doc字段里可能嵌套很多层,这里做了一个简单的拍平
Dataset<Row> hudiDf = cdcDf
.select(
col("id"),
col("op"),
from_json(col("doc").cast("string"), readUserSchema()).as("data"),
col("source_ts")
)
.select(
col("id"),
col("op"),
col("data.*"), // 展开data下的所有字段
col("source_ts")
);
// 4. 用foreachBatch把每批数据写入Hudi
// foreachBatch允许我们复用批量写的代码,对流式数据做去重或再加工
StreamingQuery query = hudiDf.writeStream()
.foreachBatch((batchDf, batchId) -> {
writeBatchToHudi(batchDf, batchId);
})
.outputMode(OutputMode.Update())
.trigger(Trigger.ProcessingTime("1 minute")) // 每分钟处理一次增量
.start();
query.awaitTermination();
}
// 定义Mongo文档的schema,按你的业务字段来写
private static StructType readUserSchema() {
return new StructType()
.add("name", DataTypes.StringType)
.add("phone", DataTypes.StringType)
.add("email", DataTypes.StringType)
.add("address", DataTypes.StringType);
}
// 这个方法把一批数据写入Hudi表
private static void writeBatchToHudi(Dataset<Row> batchDf, long batchId) {
// 先把“delete”操作过滤出去,因为Hudi的upsert操作本身不负责删除
// 如果MongoDB里发生的是删除,我们需要在Hudi里也删除对应的记录
Dataset<Row> upsertDf = batchDf.filter(col("op").notEqual("delete")).drop("op");
Dataset<Row> deleteDf = batchDf.filter(col("op").equalTo("delete")).select("id");
// 写入更新和插入
if (!upsertDf.isEmpty()) {
upsertDf.write()
.format("hudi")
.option(DataSourceWriteOptions.TABLE_TYPE_OPT_KEY(), HoodieTableType.MERGE_ON_READ.name())
.option(DataSourceWriteOptions.RECORDKEY_FIELD_OPT_KEY(), "id")
.option(DataSourceWriteOptions.PRECOMBINE_FIELD_OPT_KEY(), "source_ts")
.option(DataSourceWriteOptions.PARTITIONPATH_FIELD_OPT_KEY(), "source_ts")
.option(DataSourceWriteOptions.OPERATION_OPT_KEY(), DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL())
.option(HoodieWriteConfig.TABLE_NAME, "hudi_users")
.mode(SaveMode.Append)
.save("/warehouse/hudi/users");
}
// 删除操作走单独的路,用Hudi的删除写入API
if (!deleteDf.isEmpty()) {
deleteDf.write()
.format("hudi")
.option(DataSourceWriteOptions.TABLE_TYPE_OPT_KEY(), HoodieTableType.MERGE_ON_READ.name())
.option(DataSourceWriteOptions.RECORDKEY_FIELD_OPT_KEY(), "id")
.option(DataSourceWriteOptions.PRECOMBINE_FIELD_OPT_KEY(), "source_ts")
.option(DataSourceWriteOptions.OPERATION_OPT_KEY(), DataSourceWriteOptions.DELETE_OPERATION_OPT_VAL())
.option(HoodieWriteConfig.TABLE_NAME, "hudi_users")
.mode(SaveMode.Append)
.save("/warehouse/hudi/users");
}
}
}
上面的代码有个需要注意的地方:clusterTime从MongoDB变更流里读出来其实是BSON时间戳,不是Unix时间戳。我们需要在MongoDB连接器里配置格式转换,或者用$toLong表达式。示例里直接cast成long是简化写法,实际项目中要先把BSON时间戳转成可比较的毫秒值,否则合并时会出错。
另外,foreachBatch模式特别适合把流式数据和批量逻辑结合在一起。我们可以在同一个批次里对数据做去重、清理脏数据,然后再批量写Hudi。这样比一条一条写稳定得多。
3.4 在Hudi上做时间旅行查询
数据同步过来之后,我们最常用的一个功能是查“某个时间点之前的数据”。Hudi的COW表天然支持时间旅行,MOR表在最新版也支持了。我们可以这样查:
import org.apache.spark.sql.*;
public class TimeTravelQuery {
public static void main(String[] args) {
SparkSession spark = SparkSession.builder()
.appName("hudi-time-travel")
.config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog")
.getOrCreate();
// 直接读取Hudi表在指定时间点的快照
Dataset<Row> snapshotDf = spark.read()
.format("hudi")
.option("as.of.instant", "20250101000000") // 这个时间是Hudi的commit时间,不是业务时间
.load("/warehouse/hudi/users");
snapshotDf.createOrReplaceTempView("users_snapshot");
// 查询在某个时间点之前,手机号叫“张三”的用户
spark.sql("SELECT name, phone, email FROM users_snapshot WHERE name = '张三'").show();
}
}
这里要留意,as.of.instant用的是Hudi的提交时间,不是我们业务数据的source_ts。如果你希望按业务时间旅行,得自己用source_ts过滤。
四、迁移过程中遇到的那些坑
第一个坑:MongoDB的id不是简单的字符串,而是ObjectId。直接把_id转成字符串时,如果不做正确处理,会导致Hudi表里的主键跟MongoDB里的主键不一致。我们用StringType包了一层,但要注意ObjectId转出来的字符串可能带特殊字符,最好统一用hex表示。
第二个坑:全量同步和增量同步同时跑,会重复写同一行数据。我们的解决方法是设了source_ts,但还有更复杂的场景,比如增量作业先读到一条记录,全量作业后读到同一条记录,如果没有仔细对比precombine字段的值,就会把新数据覆盖掉。所以全量作业的时间戳必须设置成“开始导出的时间点”,而不是“每条记录自己的修改时间”。
第三个坑:MongoDB的变更流默认只会保留关于操作的一些元数据,不包含完整的文档。一定记得开fullDocument配置,否则你拿到的记录只是“更新了哪个字段”,而不是“更新之后的全量文档”。我们一开始没开,写进Hudi的很多行数据都是null。
第四个坑:Hudi的MOR表查询时,如果log文件太多,查询性能会变得很慢。我们最初太依赖MOR,忘记做compaction,结果查询Hudi表要等十几秒才出结果。后来配置了自动compaction,才恢复正常。
五、这套方案里值得注意的细节
应用场景方面,这套架构特别适合“按主键更新的OLTP库”同步历史,不光是MongoDB,MySQL、PostgreSQL也都能用同样的思路。如果你的业务数据是日志型、只追加不更新,那Hudi的另外一套玩法更适合。
技术优缺点上,Hudi的最大优点是“表和流统一”:同一份数据既能当表查,又能当流读,下游可以接实时计算,也可以直接跑离线报表。缺点也很明显:依赖Spark或Flink,运维成本不低;对HDFS或对象存储的稳定性要求高;小文件问题要经常处理。
注意事项方面,首先一定要设计好分区字段。我们一开始用时间戳分区,导致每天几十万个小文件,后来改成天级分区好多了。其次,要监控Hudi表的commit时间,如果长时间没有提交,说明同步作业挂了。再有,就是清理Hudi的历史文件版本,否则存储量会一直涨,我们每天都会执行runClustering和clean操作。
另外,删除操作不能走普通的upsert,要用Hudi的delete操作。MongoDB里的很多数据其实并不是真正删掉,而是打个软删标志,所以我们大多数时候把人家的status字段改成deleted,而不是物理删除。
六、总结
从MongoDB迁到Hudi,我们相当于给自己建了一个“时光机”。线上MongoDB只负责处理事务,历史数据全部沉淀在Hudi里,能快照、能回放、能对账。整个过程里,最大的收获不是弄明白Hudi怎么配参数,而是理解了“全量+增量+幂等合并”这套思路。只要把主键、时间戳和操作类型这三样东西控制住,不管数据源头怎么折腾,Hudi都能稳稳接住。
当然,Hudi不是银弹。如果你团队里没人熟悉Spark,也没有人愿意花时间守护流任务,那么这套方案会让你很痛苦。但如果你跟我们一样,需要做长期的数据回溯和全量历史分析,Hudi绝对值得试一把。把自己的那杯MongoDB“鲜牛奶”倒进Hudi这口“大锅”,虽然过程有点颠簸,但煮出来的味道,是真的香。
评论
围绕“将Hudi作为CDC全量历史数据存储层:从MongoDB到Hudi的迁移经验”参与讨论