Spark对比十亿级跨库数据写入Oracle性能优化求助
大数据量下Spark对比数据库表性能优化问题
我是Apache Spark新手,需要对比两个数据库中的table1和table2,把差异数据写入第三个数据库的change_aud表。小数据量(10000行)测试时功能正常,但处理各10亿行数据时,耗时超4小时仍无结果。
环境配置
- 内存:32GB
- 可用核心:8
现有代码
SparkSession spark = SparkSession.builder().master("local[*]").appName("SparkandOracledbTest").getOrCreate(); Dataset<Row> df = spark.read().format("jdbc") .option("url", "jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", "table1").option("user", "user").option("password", "password") .load(); spark.conf().set("spark.sql.sources.partitionOverwriteMode", "dynamic"); Dataset<Row> df1 = spark.read().format("jdbc") .option("url", "jdbc:oracle:thin:@jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", "table2").option("user", "user").option("password", "password") .load(); Dataset<Row> dfChange = df.except(df1); dfChange.write().partitionBy("TRACKING_ID").format("jdbc") .option("url", "jdbc:oracle:thin:@jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", "change_aud").option("user", "user").option("password", "password") .mode(SaveMode.Append).save();
优化建议
1. 分区读取JDBC数据,避免全量加载
直接全量加载10亿行数据会耗尽内存并拖慢IO,通过JDBC分区参数拆分数据为多个并行任务读取:
Dataset<Row> df = spark.read().format("jdbc") .option("url", "jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", "table1") .option("user", "user") .option("password", "password") .option("partitionColumn", "TRACKING_ID") // 选数据分布均匀的列作为分区键 .option("lowerBound", "1") // 分区列最小值 .option("upperBound", "1000000000") // 分区列最大值 .option("numPartitions", "16") // 分区数建议设为核心数的2倍,提升并行度 .load();
注意:分区列需为数值/日期类型,且数据分布均匀,避免单分区数据量过大。
2. 替换except为高效差异对比逻辑
except需要全量加载两个数据集后全局对比,10亿级数据下效率极低,推荐两种方案:
- 方案一:数据库层面计算差异
把对比逻辑下推到Oracle,利用数据库索引和查询优化能力,Spark仅读取结果:
String diffSql = "(SELECT * FROM table1 MINUS SELECT * FROM table2) diff_result"; Dataset<Row> dfChange = spark.read().format("jdbc") .option("url", "jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", diffSql) .option("user", "user") .option("password", "password") .load();
- 方案二:Spark分区局部对比
若必须用Spark对比,先按主键分区,再在每个分区内做差异对比,减少全局shuffle:
// 两个数据集按TRACKING_ID分区 Dataset<Row> partitionedDf = df.repartition("TRACKING_ID"); Dataset<Row> partitionedDf1 = df1.repartition("TRACKING_ID"); // 分区内做except Dataset<Row> dfChange = partitionedDf.except(partitionedDf1);
3. 优化JDBC写入性能
当前写入配置未充分利用并行能力,调整参数提升写入效率:
dfChange.write().format("jdbc") .option("url", "jdbc:oracle:thin:@localhost:1521/dbname") .option("dbtable", "change_aud") .option("user", "user") .option("password", "password") .option("batchSize", "10000") // 批量写入,减少数据库连接开销 .option("numPartitions", "16") // 并行写入分区数,和读取分区数匹配 .mode(SaveMode.Append).save();
注意:JDBC写入时partitionBy不生效,需通过numPartitions控制并行度。
4. 调整Spark核心配置适配硬件
根据32GB内存、8核心的配置,优化Spark资源参数:
SparkSession spark = SparkSession.builder() .master("local[8]") // 明确指定8核心,避免local[*]的不确定性 .appName("SparkandOracledbTest") .config("spark.driver.memory", "24g") // 分配24GB内存给Driver(local模式下Driver=Executor) .config("spark.sql.shuffle.partitions", "32") // 减少shuffle分区数,默认200过多 .config("spark.default.parallelism", "16") // 设置默认并行度,匹配核心数 .getOrCreate();
留2-4GB内存给系统,避免OOM。
5. 排查并解决数据倾斜
若TRACKING_ID存在大量重复值,会导致部分分区数据量过大:
- 更换数据分布更均匀的分区列,或使用复合列分区
- 对倾斜键添加随机前缀,拆分大分区为多个小分区
内容的提问来源于stack exchange,提问作者Ghaleppa Vijaykumar
相关产品推荐
相关产品推荐

