Scio Scala:如何将SCollection[Long]转换为Long以计算差值?
解决方案
问题原因
pipe.map(_.a).max返回的是SCollection[Long]类型,这是Scio的分布式集合类型,并非单个Long值,所以asInstanceOf[Long]强制类型转换必然失效——两者属于完全不同的类型体系。
两种可行处理方式
1. 分布式管道内计算差值(推荐)
保持计算在Scio管道中进行,利用cross操作将单元素的max和min集合交叉,再计算差值:
val maxCol = pipe.map(_.a).max val minCol = pipe.map(_.a).min // 结果仍为SCollection[Long],可继续在管道中处理 val diffCol: SCollection[Long] = maxCol.cross(minCol).map { case (maxVal, minVal) => maxVal - minVal }
2. 提取本地Long值(适合驱动端获取结果)
如果需要将结果拿到本地作为普通Long变量使用,可以通过run().waitForResult()获取集合中的单元素:
// 阻塞等待作业完成,获取max的本地值 val maxVal: Long = pipe.map(_.a).max.run().waitForResult().head // 同理获取min的本地值 val minVal: Long = pipe.map(_.a).min.run().waitForResult().head // 直接计算差值 val diff: Long = maxVal - minVal
注意:如果原始数据可能为空,head会抛出异常,建议用headOption做安全处理:
val maxOpt: Option[Long] = pipe.map(_.a).max.run().waitForResult().headOption val minOpt: Option[Long] = pipe.map(_.a).min.run().waitForResult().headOption // 仅当max和min都存在时计算差值 val diffOpt: Option[Long] = for { m <- maxOpt n <- minOpt } yield m - n
内容的提问来源于stack exchange,提问作者Lara
相关产品推荐
相关产品推荐

