You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Scalatest中对比相同Spark Dataset断言失败问题

问题分析与解决方案

核心问题1:数据类型隐式转换不匹配

你的driveAggMonthly类中numberOfDrives字段是BigInt类型,但nr_drives_by_month函数里,count(col("driveId"))返回的是Long类型。虽然Spark的Schema显示两者都是bigint,但在Scala层面,Long和BigInt是不同的类型,这会导致Dataset的实际数据类型不兼容,进而触发断言失败。

核心问题2:Dataset直接用ScalaTest断言不可靠

ScalaTest的===操作符无法正确对比分布式的Dataset对象,即便数据和Schema看起来一致,底层的Dataset实例元数据(比如分区、执行计划)差异也会导致断言失败,必须将数据收集到本地集合后再对比。


修复步骤

1. 修正聚合函数的类型转换

修改nr_drives_by_month函数,将count的结果显式转换为BigInt:

def nr_drives_by_month(input_ds: Dataset[Drive]): Dataset[driveAggMonthly] = {
  input_ds
    .groupBy(month(to_date(col("date"))).as("Month"))
    // 将count的Long结果转为BigInt
    .agg(count(col("driveId")).cast(BigIntType).as("numberOfDrives"))
    .as[driveAggMonthly]
}

或者在Scala层面做类型转换:

.agg(count(col("driveId")).as[Long].map(BigInt(_)).as("numberOfDrives"))

2. 正确的断言方式

不要直接对比Dataset对象,而是将两者的数据收集到本地Seq后再断言:

// 收集数据到本地集合
val expectedData = expected_ds_t1.collect().toSeq
val resultData = result_ds_t1.collect().toSeq

// 断言数据一致
assert(resultData === expectedData)

如果需要额外验证Schema一致性,可以添加:

assert(result_ds_t1.schema === expected_ds_t1.schema)

验证说明

  • 类型修正后,Dataset中的numberOfDrives字段会真正匹配BigInt类型,避免Scala层面的类型不兼容问题。
  • 收集到本地后对比是分布式数据测试的标准做法,能绕过Dataset实例本身的元数据差异,只对比实际业务数据。

内容的提问来源于stack exchange,提问作者Vanessa Hydoski

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 14:56:06