如何使用Spark将相同DataFrame写入含源数据源的多目标?
Spark 2.4.4 解决读写同数据源后的重复数据问题
针对你遇到的「从单一数据源读取数据,先写回原数据源再写入其他数据源时出现重复数据」的问题,结合Spark 2.4.4和Java 8环境,以下是仅通过Spark实现的可行方案:
方案一:Checkpoint 打破血缘依赖
Checkpoint会将DataFrame的快照写入可靠分布式存储(如HDFS),同时切断RDD/DataFrame的血缘关系,彻底避免第二次写入时重新读取已被修改的原数据源。
操作步骤:
- 配置Checkpoint目录(需使用分布式存储,本地目录仅适合单节点测试):
spark.sparkContext().setCheckpointDir("hdfs://your-cluster/checkpoint-dir"); - 读取原始数据后执行Checkpoint并强制物化:
Dataset<Row> rawDf = spark.read().format("mongodb").load(); // 示例读取Mongo数据源 // 执行eager checkpoint(Spark 2.4中需触发action完成物化) Dataset<Row> checkpointedDf = rawDf.checkpoint(); checkpointedDf.count(); // 触发action,将数据写入checkpoint目录并切断血缘 - 使用checkpointedDf依次完成两次写入:
// 写回原数据源 checkpointedDf.write().mode(SaveMode.Overwrite).format("mongodb").save(); // 写入第二个数据源 checkpointedDf.write().mode(SaveMode.Append).format("parquet").save("hdfs://your-cluster/target-data");
方案二:外部中间存储落地快照数据
将原始读取的数据先写入一个临时外部存储(如HDFS Parquet、临时Mongo集合),后续两次写入均基于这个快照数据,完全隔离原数据源的修改影响。
操作步骤:
- 读取原始数据并写入临时存储:
Dataset<Row> rawDf = spark.read().format("csv").load("hdfs://your-cluster/source-data"); // 写入临时Parquet存储(可替换为其他支持的格式) rawDf.write().mode(SaveMode.Overwrite).parquet("hdfs://your-cluster/temp-snapshot"); - 从临时存储读取快照数据:
Dataset<Row> snapshotDf = spark.read().parquet("hdfs://your-cluster/temp-snapshot"); - 用snapshotDf完成两次写入:
// 写回原CSV数据源 snapshotDf.write().mode(SaveMode.Overwrite).csv("hdfs://your-cluster/source-data"); // 写入MongoDB数据源 snapshotDf.write().mode(SaveMode.Append).format("mongodb").save();
注:临时存储可在任务完成后删除,避免占用空间。
方案三:规避SPARK-24596的持久化方式
针对SPARK-24596导致的磁盘persist缓存失效问题,可改用内存优先的持久化级别,并在写入前强制完成数据物化:
操作步骤:
- 读取数据后配置持久化级别并强制加载:
Dataset<Row> rawDf = spark.read().format("kafka").load(); // 使用MEMORY_AND_DISK级别,避免纯磁盘persist的失效问题 rawDf.persist(StorageLevel.MEMORY_AND_DISK()); rawDf.count(); // 触发action,将数据完全加载到缓存 - 直接使用缓存后的DataFrame完成两次写入:
rawDf.write().mode(SaveMode.Overwrite).format("kafka").save(); rawDf.write().mode(SaveMode.Append).format("parquet").save("hdfs://your-cluster/target-data");
注:该方案仅适合数据量能适配集群内存的场景,大数据量下优先选择前两种方案。
方案对比
| 方案 | 适用场景 | 优势 | 局限性 |
|---|---|---|---|
| Checkpoint | 分布式集群、大数据量场景 | 无需额外管理临时存储,自动切断血缘 | 依赖分布式存储,单节点测试需配置本地目录 |
| 外部中间存储 | 全场景兼容,尤其是异构数据源 | 逻辑清晰,完全隔离原数据源修改影响 | 需手动清理临时存储 |
| 内存优先持久化 | 小数据量、内存充足的场景 | 操作简单,无需额外存储开销 | 受内存容量限制,存在OOM风险 |
内容的提问来源于stack exchange,提问作者maxime G
相关产品推荐
相关产品推荐

