Spark Scala左连接未达预期,咨询逻辑正确性与输出结果
问题描述
数据结构定义
case class ReconEntity(rowId: String, groupId: String, amounts: List[Amount], processingDate: Long, attributes: Map[String, String], entityType: String, isDuplicate: String)
输入数据集
第一个Dataset(inputEntries)
+-----+-------+------------------+--------------+----------+-----------+ |rowId|groupId| amount|processingDate|attributes|isDuplicate| +-----+-------+------------------+--------------+----------+-----------+ | C1| G1|USD,10.00000000...| 1551021334| rowId,C1| false| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C2| G2|USD,2.000000000...| 1551011017| rowId,C2| false| | C3| G2|USD,6.000000000...| 1551011459| rowId,C3| false| | C3| G2|USD,6.000000000...| 1551011017| rowId,C3| true| +-----+-------+------------------+--------------+----------+-----------+
第二个Dataset(otherEntries)
+-----+-------+------------------+--------------+----------+-----------+ |rowId|groupId| amount|processingDate|attributes|isDuplicate| +-----+-------+------------------+--------------+----------+-----------+ | C2| G2|USD,2.000000000...| 1551011017| rowId,C2| false| | C3| G2|USD,6.000000000...| 1551011459| rowId,C3| false| +-----+-------+------------------+--------------+----------+-----------+
期望结果
+-----+-------+------------------+--------------+----------+-----------+ |rowId|groupId| amount|processingDate|attributes|isDuplicate| +-----+-------+------------------+--------------+----------+-----------+ | C1| G1|USD,10.00000000...| 1551021334| rowId,C1| true| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C1| G1|USD,10.00000000...| 1551011017| rowId,C1| true| | C2| G2|USD,2.000000000...| 1551011017| rowId,C2| false| | C3| G2|USD,6.000000000...| 1551011459| rowId,C3| false| | C3| G2|USD,6.000000000...| 1551011017| rowId,C3| true| +-----+-------+------------------+--------------+----------+-----------+
当前实现代码
inputEntries.as("inputDataset").join(otherEntries.as("otherDataset"), col("inputDataset.rowId") === col("otherDataset.rowId"), "left") .select( col("inputDataset.rowId"), col("inputDataset.groupId"), col("inputDataset.amounts"), col("inputDataset.processingDate"), col("inputDataset.attributes"), col("inputDataset.entityType"), when( col("otherDataset.rowId").isNull, TRUE ).otherwise(col("inputDataset.isDuplicate")).as(IS_DUPLICATE) ).as[ReconEntity]
本地运行正常,但提交Spark任务后结果不符合预期,需确认逻辑正确性及左连接的实际输出。
问题分析与修正
原逻辑的核心问题
你的左连接逻辑本身方向是对的,但存在类型不匹配的隐患:ReconEntity中isDuplicate是String类型,代码里却直接用布尔值TRUE赋值。本地环境可能隐式处理了类型转换,但Spark集群环境的类型校验更严格,会导致赋值错误,这是集群结果不符合预期的主要原因。
左连接的实际中间输出
以你的数据为例,左连接后的关键字段结果如下:
| inputDataset.rowId | otherDataset.rowId | inputDataset.isDuplicate |
|---|---|---|
| C1 | null | false |
| C1 | null | true |
| C1 | null | true |
| C1 | null | true |
| C2 | C2 | false |
| C3 | C3 | false |
| C3 | C3 | true |
可以看到:
- C1的所有记录在otherDataset中无匹配,
otherDataset.rowId为null - C2、C3的记录都能匹配到otherDataset中的对应rowId,
otherDataset.rowId不为null
修正后的代码
将TRUE改为字符串类型的"true",确保和字段定义一致:
inputEntries.as("inputDataset").join(otherEntries.as("otherDataset"), col("inputDataset.rowId") === col("otherDataset.rowId"), "left") .select( col("inputDataset.rowId"), col("inputDataset.groupId"), col("inputDataset.amounts"), col("inputDataset.processingDate"), col("inputDataset.attributes"), col("inputDataset.entityType"), when( col("otherDataset.rowId").isNull, lit("true") ).otherwise(col("inputDataset.isDuplicate")).as(IS_DUPLICATE) ).as[ReconEntity]
更高效的优化方案
如果只需要判断rowId是否存在于otherDataset中,没必要做全量左连接,先提取去重的rowId集合再判断,性能更优:
// 提取otherDataset中所有存在的rowId并去重 val existingRowIds = otherEntries.select("rowId").distinct().as[String] // 直接判断rowId是否存在,避免全量连接 inputEntries.withColumn(IS_DUPLICATE, when(not(exists(existingRowIds, col("rowId") === col("inputDataset.rowId"))), lit("true")) .otherwise(col("isDuplicate")) ).as[ReconEntity]
内容的提问来源于stack exchange,提问作者anonymous
相关产品推荐
相关产品推荐

