【问题标题】:filter with instanceOf Tuple使用 instanceOf 元组过滤
【发布时间】:2018-07-26 19:40:14
【问题描述】:

我试图找到单词的共现。以下是我正在使用的代码。

val dataset = df.select("entity").rdd.map(row => row.getList(0)).filter(r => r.size() > 0).distinct()
println("dataset")

dataset.take(10).foreach(println)

示例数据集

dataset
[aa]
[bb]
[cc]
[dd]
[ee]
[ab, ac, ad]
[ff]
[ef, fg]
[ab, gg, hh]

代码片段

case class tupleIn(a: String,b: String)
case class tupleOut(i: tupleIn, c: Long)
val cooccurMapping = dataset.flatMap(
list => {
    list.toArray().map(e => e.asInstanceOf[String].toLowerCase).flatMap(
        ele1 => {
            list.toArray().map(e => e.asInstanceOf[String].toLowerCase).map(ele2 => {
                if (ele1 != ele2) {
                    ((ele1, ele2), 1L)
                }
            })
        })
})

如何从中过滤?

我试过了

.filter(e => e.isInstanceOf[Tuple2[(String, String), Long]])

:121: 警告:无果类型测试:Unit 类型的值也不能是 ((String, String), Long) .filter(e => e.isInstanceOf[Tuple2[(String, String), Long]]) ^

:121: 错误:isInstanceOf 无法测试值类型是否为引用。 .filter(e => e.isInstanceOf[Tuple2[(String, String), Long]])

.filter(e => e.isInstanceOf[tupleOut])

:122: 警告:无果类型测试:单元类型的值 也不能是 coocrTupleOut .filter(e => e.isInstanceOf[tupleOut]) ^ :122: 错误:isInstanceOf 无法测试值类型是否为引用。 .filter(e => e.isInstanceOf[tupleOut])

如果我映射

.map(e => e.asInstanceOf[Tuple2[(String, String), Long]])

上面的 sn-p 工作正常,但一段时间后会出现此异常:

java.lang.ClassCastException: scala.runtime.BoxedUnit 不能被强制转换 到 scala.Tuple2 在 $line84834447093.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw $$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$anonfun$2$$anonfun$9.apply( :123) 在 $line84834447093.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw $$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$anonfun$2$$anonfun$9.apply( :123) 在 scala.collection.Iterator$$anon$11.next(Iterator.scala:409) 在 scala.collection.Iterator$$anon$13.hasNext(Iterator.scala:462) 在 org.apache.spark.util.collection.ExternalSorter.insertAll(ExternalSorter.scala:191) 在 org.apache.spark.shuffle.sort.SortShuffleWriter.write(SortShuffleWriter.scala:63) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:96) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:53) 在 org.apache.spark.scheduler.Task.run(Task.scala:108) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:338) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) 在 java.lang.Thread.run(Thread.java:748)

为什么instanceOf 不在filter() 工作,但在map() 工作

【问题讨论】:

标签: scala apache-spark


【解决方案1】:

您的代码的结果是 Unit 类型的项目的集合,因此过滤器和映射都不会迭代(请注意,映射 as 因此它会转换为您想要的类型 wheres as is 检查类型)

无论如何,如果我正确理解您的意图,您可以使用 spark 的内置函数得到您想要的:

val l=List(List("aa"),List("bb","vv"),List("bbb"))
val rdd=sc.parallelize(l)
val df=spark.createDataFrame(rdd,"data")

import org.apache.spark.sql.functions._
val ndf=df.withColumn("data",explode($"data"))
val cm=ndf.select($"data".as("elec1")).crossJoin(ndf.select($"data".as("elec2"))).withColumn("cnt",lit(1L))
val coocurenceMap=cm.filter($"elec1" !== $"elec2")

【讨论】:

  • 谢谢。但我试图在不使用 spark 的数据帧 API 的情况下实现这一目标。
  • 正如我告诉您的那样,您的地图和过滤器都会获取 Unit 列表(您在异常中看到的 BoxedUnit)而不是元组。此外,当你让它工作时,它的效率会低于上面的代码
猜你喜欢
  • 2018-12-25
  • 2019-01-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-08-09
  • 2016-09-11
  • 2021-03-21
  • 1970-01-01
相关资源
最近更新 更多