Erro do compilador de renderização Scastie como “value countByValue não é membro de org.apache.spark.sql.Dataset [String]”

Sep 10 2020

Olá, estou tentando encontrar o histograma de avaliações usando o programa scastie ... aqui está a implementação

configurações do sbet no scastie

        scalacOptions ++= Seq(
          "-deprecation",
          "-encoding", "UTF-8",
          "-feature",
          "-unchecked"
        )

            libraryDependencies ++= Seq(
              "org.apache.spark" %% "spark-core" % "2.4.3",
              "org.apache.spark" %% "spark-sql" % "2.4.3"
            )

código real em scastie

                    import org.apache.spark.sql.SparkSession
                    import org.apache.spark._
                    import org.apache.spark.SparkContext._
                    import org.apache.spark.sql.SparkSession
                    import org.apache.log4j._


                        object TestApp extends App {
                      lazy implicit val spark = 
                      SparkSession.builder().master("local").appName("spark_test").getOrCreate()
                      
                      import spark.implicits._ // Required to call the .toDF function later
                      
                      val html = scala.io.Source.fromURL("http://files.grouplens.org/datasets/movielens/ml- 
     
                      100k/u.data").mkString // Get all rows as one string
                      val seqOfRecords = html.split("\n") // Split based on the newline characters
                                     .filter(_ != "") // Filter out any empty lines
                                     .toSeq // Convert to Seq so we can convert to DF later
                                     .map(row => row.split("\t")) 
                                     .map { case Array(f1,f2,f3,f4) => (f1,f2,f3,f4) } 
                      
                      val df = seqOfRecords.toDF("col1", "col2", "col3", "col4") 
                      
                      val ratings = df.map(x => x.toString().split("\t")(2))
                      
                      

                    // Count up how many times each value (rating) occurs
                    val results = ratings.countByValue()

                    // Sort the resulting map of (rating, count) tuples
                    val sortedResults = results.toSeq.sortBy(_._1)

                    // Print each result on its own line.
                    sortedResults.foreach(println)

                      spark.close() 
                    }

Erro ao entrar no scastie

valor countByValue não é membro de org.apache.spark.sql.Dataset [String]

alguém pode ajudar na limpeza

================================================= Código revisado apresentando erro diferente no Scastie agora

                    java.lang.ExceptionInInitializerError
                        at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
                        at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:127)
                        at org.apache.spark.sql.execution.TakeOrderedAndProjectExec.executeCollect(limit.scala:136)
                        at org.apache.spark.sql.Dataset.org$apache$spark$sql$Dataset$$collectFromPlan(Dataset.scala:3383) at org.apache.spark.sql.Dataset$$anonfun$head$1.apply(Dataset.scala:2544)
                        at org.apache.spark.sql.Dataset$$anonfun$head$1.apply(Dataset.scala:2544) at org.apache.spark.sql.Dataset$$anonfun$53.apply(Dataset.scala:3364) at org.apache.spark.sql.execution.SQLExecution$$anonfun$withNewExecutionId$1.apply(SQLExecution.scala:78)
                        at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:125) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:73)
                        at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3363)
                        at org.apache.spark.sql.Dataset.head(Dataset.scala:2544)
                        at org.apache.spark.sql.Dataset.take(Dataset.scala:2758)
                        at org.apache.spark.sql.Dataset.getRows(Dataset.scala:254)
                        at org.apache.spark.sql.Dataset.showString(Dataset.scala:291)
                        at org.apache.spark.sql.Dataset.show(Dataset.scala:745)
                        at org.apache.spark.sql.Dataset.show(Dataset.scala:704)
                        at org.apache.spark.sql.Dataset.show(Dataset.scala:713)
                        at TestApp$.delayedEndpoint$TestApp$1(main.scala:22) at TestApp$delayedInit$body.apply(main.scala:4) at scala.Function0$class.apply$mcV$sp(Function0.scala:34)
                        at scala.runtime.AbstractFunction0.apply$mcV$sp(AbstractFunction0.scala:12)
                        at scala.App$$anonfun$main$1.apply(App.scala:76) at scala.App$$anonfun$main$1.apply(App.scala:76)
                        at scala.collection.immutable.List.foreach(List.scala:392)
                        at scala.collection.generic.TraversableForwarder$class.foreach(TraversableForwarder.scala:35) at scala.App$class.main(App.scala:76)
                        at TestApp$.main(main.scala:4) at TestApp.main(main.scala) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at sbt.Run.invokeMain(Run.scala:115) at sbt.Run.execute$1(Run.scala:79)
                        at sbt.Run.$anonfun$runWithLoader$4(Run.scala:92) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at sbt.util.InterfaceUtil$$anon$1.get(InterfaceUtil.scala:10) at sbt.TrapExit$App.run(TrapExit.scala:257)
                        at java.lang.Thread.run(Thread.java:748)
                    Caused by: com.fasterxml.jackson.databind.JsonMappingException: Incompatible Jackson version: 2.9.8
                        at com.fasterxml.jackson.module.scala.JacksonModule$class.setupModule(JacksonModule.scala:64) at com.fasterxml.jackson.module.scala.DefaultScalaModule.setupModule(DefaultScalaModule.scala:19) at com.fasterxml.jackson.databind.ObjectMapper.registerModule(ObjectMapper.java:751) at org.apache.spark.rdd.RDDOperationScope$.<init>(RDDOperationScope.scala:82)
                        at org.apache.spark.rdd.RDDOperationScope$.<clinit>(RDDOperationScope.scala)
                        ... 40 more

aqui está o código atualizado em scastie

                import org.apache.spark.sql.SparkSession
                import org.apache.spark.sql.functions.col

                object TestApp extends App {
                  lazy implicit val spark = SparkSession.builder().master("local").appName("spark_test").getOrCreate()
                  
                  import spark.implicits._ // Required to call the .toDF function later
                  
                  val html = scala.io.Source.fromURL("http://files.grouplens.org/datasets/movielens/ml-100k/u.data").mkString // Get all rows as one string
                  val seqOfRecords = html.split("\n") // Split based on the newline characters
                                 .filter(_ != "") // Filter out any empty lines
                                 .toSeq // Convert to Seq so we can convert to DF later
                                 .map(row => row.split("\t")) // Split each line on tab character to make an Array of 4 String each
                                 .map { case Array(f1,f2,f3,f4) => (f1,f2,f3,f4) } // Convert that Array[String] into Array[(String, String, String, String)] 
                  
                  val df = seqOfRecords.toDF("col1", "col2", "col3", "col4") // Give whatever column names you want
                  
                  df.select("col3").groupBy("col3").count.sort(col("count").desc).show()

                  spark.close() // don't forget to close(), otherwise scastie won't let you create another session so soon.
                }

Respostas

1 kfkhalili Sep 11 2020 at 00:18

Primeira parte da sua pergunta: Portanto, o principal problema em seu código é a tentativa de dividir por tabulação \t. Seus registros não contêm guias, como expliquei em meu comentário.

O fato é que, quando você mapeia por meio do df, está acessando cada org.apache.spark.sql.Rowobjeto, por exemplo, df.firsté [196,242,3,881250949]. Você pode transformar isso em a String, mas não há \t(caractere de tabulação) para dividir, então ele simplesmente retornará um Stringcomo está em Array[String]com apenas um elemento, portanto, acessar o segundo elemento retorna um java.lang.ArrayIndexOutOfBoundsException.

Aqui está uma demonstração:

// We get the first row and brute force convert it toString()
df.head.toString
//res21: String = [196,242,3,881250949] <- See? No tab anywhere

df.head.toString.split("\t")
//res22: Array[String] = Array([196,242,3,881250949]) <- Returns the string as is in an Array

res22(0)
//res24: String = [196,242,3,881250949] <- First Element

res22(1)
//java.lang.ArrayIndexOutOfBoundsException: 1 <- No second (or third) element found, hence the "out of bounds" exception.
//  ... 55 elided

Eu entendi pelo seu comentário que você está tentando obter a terceira coluna. A beleza de usar um DataFrameé que você pode simplesmente selectdefinir a coluna desejada pelo nome. Você pode então groupBy(isso retorna um RelationalGroupedDataset ) e usar o countmétodo para agregar.

import org.apache.spark.sql.functions.col
df.select("col3").groupBy("col3").count.sort(col("count").desc).show()
//+----+-----+
//|col3|count|
//+----+-----+
//|   4|34174|
//|   3|27145|
//|   5|21201|
//|   2|11370|
//|   1| 6110|
//+----+-----+

Segunda parte da sua pergunta: Parece que cargas Scastie uma versão mais recente com.fasterxml.jackson.core:jackson-databinddo que o que faísca 2.4.3 usos, por isso, enquanto Scastie parece versão uso 2.9.6, o Spark 2.4.3 usa uma versão antiga: 2.6.7.

A única maneira de fazer funcionar era usar uma versão mais recente do Spark e do Scala. Usos do Spark 3.0.1 2.10.0.

Em Configurações de compilação:

  • Defina Scala Versioncomo 2.12.10.
  • Definir dependências de biblioteca de configuração extra Sbt:
libraryDependencies ++= Seq(
  "org.apache.spark" %% "spark-core" % "3.0.1",
  "org.apache.spark" %% "spark-sql" % "3.0.1"
)

Ele não funciona tão bem, o navegador trava e às vezes atinge o tempo limite. Acho que Scastie ainda não está otimizado para esta versão.

Edit: Na verdade, depois que silenciei o registro, ele funciona muito melhor agora !

MAS ainda ... Você realmente deve instalar o Spark em seu computador local .

1 rich_morton Sep 10 2020 at 11:35

No momento em que você chega à ratingsvariável, está trabalhando com uma estrutura do Spark chamada Dataset. Você pode consultar a documentação que descreve o que pode e o que não pode ser feito aqui . Ele não tem um método chamado countByValuee é por isso que você obtém o erro que está vendo.

Tudo o que você tem faz sentido até chegar a esta linha:

val ratings = df.map(x => x.toString().split("\t")(2))

No momento, isso irá gerar um erro.

Se você voltar para a dfvariável, terá uma tabela que se parecerá com esta:

+----+----+----+---------+
|col1|col2|col3|     col4|
+----+----+----+---------+
| 196| 242|   3|881250949|
| 186| 302|   3|891717742|
|  22| 377|   1|878887116|
| 244|  51|   2|880606923|
| 166| 346|   1|886397596|
+----+----+----+---------+
                  

Você pode executar o comando df.show()para ver uma amostra do que está no conjunto de dados. A partir daí, acho que você está querendo uma operação que se pareça um pouco com groupBy. Dê uma olhada em alguns exemplos para ver o que fazer a seguir.