【发布时间】:2020-04-09 15:24:59
【问题描述】:
我正在研究 Spark SQL,我需要找出两个大型 CSV 之间的差异。
Diff 应该给出:-
插入的行或新记录 // 仅比较 Id 的
更改的行(不包括插入的行)- 比较所有列值
- 已删除的行 // 仅比较 Id 的
Spark 2.4.4 + Java
我正在使用 Databricks 读取/写入 CSV
Dataset<Row> insertedDf = newDf_temp.join(oldDf_temp,oldDf_temp.col(key)
.equalTo(newDf_temp.col(key)),"left_anti");
Long insertedCount = insertedDf.count();
logger.info("Inserted File Count == "+insertedCount);
Dataset<Row> deletedDf = oldDf_temp.join(newDf_temp,oldDf_temp.col(key)
.equalTo(newDf_temp.col(key)),"left_anti")
.select(oldDf_temp.col(key));
Long deletedCount = deletedDf.count();
logger.info("deleted File Count == "+deletedCount);
Dataset<Row> changedDf = newDf_temp.exceptAll(oldDf_temp); // This gives rows (New +changed Records)
Dataset<Row> changedDfTemp = changedDf.join(insertedDf, changedDf.col(key)
.equalTo(insertedDf.col(key)),"left_anti"); // This gives only changed record
Long changedCount = changedDfTemp.count();
logger.info("Changed File Count == "+changedCount);
这适用于最多 50 列左右的 CSV。
The Above code fails for one row in CSV with 300+columns, so I am sure this is not file Size problem.
但如果我有一个包含 300 多列的 CSV,那么它会因异常而失败
达到批量解决的最大迭代次数 (100) - Spark 错误
If I set the below property in Spark, It Works!!!
sparkConf.set("spark.sql.optimizer.maxIterations", "500");
但我的问题是为什么我必须设置这个?
我做错了什么吗? 或者对于具有大列的 CSV,这种行为是预期的。
我能否以任何方式对其进行优化以处理大列 CSV。
【问题讨论】:
标签: apache-spark apache-spark-sql data-science