使用Java与Spark对比DataFrame并高亮差异字段的实现方案
实现方案
要实现字段级的差异定位,核心是先通过唯一主键(本例为emp_id)关联两个DataFrame的同一条数据,逐字段对比值的差异,对差异字段的内容添加高亮标识后输出即可。
完整实现代码
import org.apache.spark.sql.Column; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import java.util.ArrayList; import java.util.List; import static org.apache.spark.sql.functions.*; public class DataFrameDiff { public static void main(String[] args) { SparkSession spark = SparkSession.builder() .appName("DataFrameDiff") .master("local[*]") .getOrCreate(); // 读取原始数据,注意要开启header读取表头 Dataset<Row> df1 = spark.read().option("header", "true").csv("/Users/dataframeOne.csv"); Dataset<Row> df2 = spark.read().option("header", "true").csv("/Users/dataframeTwo.csv"); // 配置主键和高亮规则 String primaryKey = "emp_id"; String[] allColumns = df1.columns(); // 控制台ANSI黄色高亮编码,其他场景可替换为对应格式:比如网页场景替换为<span style="color: gold;">和</span> String highlightPrefix = "\033[33m"; String highlightSuffix = "\033[0m"; // 给两个DataFrame的非主键字段加前缀,避免join后字段重名 List<Column> df1SelectCols = new ArrayList<>(); List<Column> df2SelectCols = new ArrayList<>(); List<String> compareCols = new ArrayList<>(); for (String col : allColumns) { if (col.equals(primaryKey)) { df1SelectCols.add(col(col)); df2SelectCols.add(col(col)); } else { df1SelectCols.add(col(col).alias("df1_" + col)); df2SelectCols.add(col(col).alias("df2_" + col)); compareCols.add(col); } } Dataset<Row> df1Rename = df1.select(df1SelectCols.toArray(new Column[0])); Dataset<Row> df2Rename = df2.select(df2SelectCols.toArray(new Column[0])); // 按主键关联两个DataFrame Dataset<Row> joinDf = df1Rename.join(df2Rename, primaryKey); // 构造结果列:差异字段添加高亮 List<Column> resultCols = new ArrayList<>(); resultCols.add(col(primaryKey)); Column diffFilter = lit(false); for (String col : compareCols) { // 保留修改前的原值列,不需要可以删掉 resultCols.add(col("df1_" + col).alias("old_" + col)); // 新值列:差异值自动加高亮 Column diffCol = when( // 用nvl兼容空值对比,避免空值判断失效 nvl(col("df1_" + col), lit("")).notEqual(nvl(col("df2_" + col), lit(""))), concat(lit(highlightPrefix), col("df2_" + col), lit(highlightSuffix)) ).otherwise(col("df2_" + col)) .alias("new_" + col); resultCols.add(diffCol); // 可选:添加字段差异标记列 Column isDiff = nvl(col("df1_" + col), lit("")).notEqual(nvl(col("df2_" + col), lit(""))); resultCols.add(isDiff.alias(col + "_is_diff")); // 构造过滤条件:只要有一个字段不同就保留该行 diffFilter = diffFilter.or(isDiff); } // 过滤出有差异的行并输出,show参数设为false避免转义字符被截断 Dataset<Row> diffResult = joinDf.select(resultCols.toArray(new Column[0])).filter(diffFilter); diffResult.show(false); } }
补充说明
- 代码运行后,示例中
new_emp_name列的romino会自动以黄色高亮显示,同时可查看修改前的原值和差异标记。如果不需要保留原值、差异标记列,直接删除resultCols中对应的添加逻辑即可。 - 若两个DataFrame没有唯一主键,可先调用
withColumn("row_id", monotonically_increasing_id())生成自增行号作为关联键,按行顺序对比即可。
内容的提问来源于stack exchange,提问作者JoeyOC
相关产品推荐
相关产品推荐

