小编eli*_*sah的帖子

Registring Kryo课程不起作用

我有以下代码:

val conf = new SparkConf().setAppName("MyApp")
val sc = new SparkContext(conf)
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
new conf.registerKryoClasses(new Class<?>[]{
        Class.forName("org.apache.hadoop.io.LongWritable"),
        Class.forName("org.apache.hadoop.io.Text")
    });
Run Code Online (Sandbox Code Playgroud)

但我碰到了以下错误:

')' expected but '[' found.
[error]                 new conf.registerKryoClasses(new Class<?>[]{
Run Code Online (Sandbox Code Playgroud)

我怎么解决这个问题 ?

scala kryo apache-spark

1
推荐指数
1
解决办法
293
查看次数

如何在 Python 中使用 elasticsearch 检索 1M 文档?

如何从 python 中获得 100000 个寄存器在 elasticsearch 中?MatchAll 查询仅检索 10000。

python elasticsearch elasticsearch-5

1
推荐指数
1
解决办法
4460
查看次数

Spark sqlContext全选

我使用Spark SQLContext读取数据并将其存储在变量中:

 val somevar = sqlContext.read.parquet(some_file.parquet)
Run Code Online (Sandbox Code Playgroud)

然后我希望使用select选择所有值,例如:

  somevar.select(*)
Run Code Online (Sandbox Code Playgroud)

但这不起作用.

相当于:

somevar.registerTempTable("sometable")

sqlContext.sql("SELECT * FROM sometable")
Run Code Online (Sandbox Code Playgroud)

但我不想做以前的事情.

亲切的问候.

scala apache-spark apache-spark-sql

1
推荐指数
1
解决办法
1709
查看次数

scala val _x和val x_是什么意思

我对下面的两个声明感到困惑:

val _x
val x_
Run Code Online (Sandbox Code Playgroud)

它是什么意思以及为什么我们有时使用下划线来声明变量和函数?

scala

1
推荐指数
1
解决办法
333
查看次数

从spark写入elasticsearch非常慢

我正在处理一个文本文件,并将转换后的行从Spark应用程序写入弹性搜索

input.write.format("org.elasticsearch.spark.sql")
      .mode(SaveMode.Append)
      .option("es.resource", "{date}/" + dir).save()
Run Code Online (Sandbox Code Playgroud)

这运行速度非常慢,大约需要8分钟才能写入287.9 MB/1513789条记录. 在此输入图像描述

如果网络延迟始终存在,我如何调整spark和elasticsearch设置以使其更快.

我在本地模式下使用spark,有16个内核和64GB RAM.我的elasticsearch集群有一个主节点和3个数据节点,每个节点有16个核心和64GB.

我正在阅读如下文本文件

 val readOptions: Map[String, String] = Map("ignoreLeadingWhiteSpace" -> "true",
  "ignoreTrailingWhiteSpace" -> "true",
  "inferSchema" -> "false",
  "header" -> "false",
  "delimiter" -> "\t",
  "comment" -> "#",
  "mode" -> "PERMISSIVE")
Run Code Online (Sandbox Code Playgroud)

....

val input = sqlContext.read.options(readOptions).csv(inputFile.getAbsolutePath)
Run Code Online (Sandbox Code Playgroud)

elasticsearch apache-spark elasticsearch-5 elasticsearch-spark

1
推荐指数
1
解决办法
1305
查看次数

spark-sql内置的dayofmonth函数返回奇怪的结果

由于某些奇怪的原因,dayofmonthspark中的函数多年来似乎返回了奇怪的值1500 or less

以下是获得的结果->

scala> spark.sql("SELECT dayofmonth('1501-02-14') ").show()
+------------------------------------+
|dayofmonth(CAST(1501-02-14 AS DATE))|
+------------------------------------+
|                                  14|
+------------------------------------+


scala> spark.sql("SELECT dayofmonth('1500-02-14') ").show()
+------------------------------------+
|dayofmonth(CAST(1500-02-14 AS DATE))|
+------------------------------------+
|                                  13|
+------------------------------------+


scala> spark.sql("SELECT dayofmonth('1400-02-14') ").show()
+------------------------------------+
|dayofmonth(CAST(1400-02-14 AS DATE))|
+------------------------------------+
|                                  12|
+------------------------------------+
Run Code Online (Sandbox Code Playgroud)

谁能解释,为什么星火如此行事?

scala apache-spark apache-spark-sql

1
推荐指数
1
解决办法
66
查看次数

如何计算多个浮点列的累积总和?

我在数据框中有 100 个按日期排序的浮动列。

ID   Date         C1       C2 ....... C100
1     02/06/2019   32.09  45.06         99
1     02/04/2019   32.09  45.06         99
2     02/03/2019   32.09  45.06         99
2     05/07/2019   32.09  45.06         99
Run Code Online (Sandbox Code Playgroud)

我需要根据 ID 和日期在累积总和中获得 C1 到 C100。

目标数据框应如下所示:

ID   Date         C1       C2 ....... C100
1     02/04/2019   32.09  45.06         99
1     02/06/2019   64.18  90.12         198
2     02/03/2019   32.09  45.06         99
2     05/07/2019   64.18  90.12         198
Run Code Online (Sandbox Code Playgroud)

我想在不从 C1-C100 循环的情况下实现这一点。

一列的初始代码:

var DF1 =  DF.withColumn("CumSum_c1", sum("C1").over(
         Window.partitionBy("ID")
        .orderBy(col("date").asc)))
Run Code Online (Sandbox Code Playgroud)

我在这里发现了一个类似的问题,但他手动为两列做了这个:Spark 中的累积总和

scala apache-spark apache-spark-sql

1
推荐指数
1
解决办法
130
查看次数

在Scala中将DataFrame转换为RDD [Map]

我想转换一个创建的数组,如:

case class Student(name: String, age: Int)
val dataFrame: DataFrame = sql.createDataFrame(sql.sparkContext.parallelize(List(Student("Torcuato", 27), Student("Rosalinda", 34))))
Run Code Online (Sandbox Code Playgroud)

当我从DataFrame收集结果时,生成的数组是一个 Array[org.apache.spark.sql.Row] = Array([Torcuato,27], [Rosalinda,34])

我正在研究在RDD [Map]中转换DataFrame,例如:

Map("name" -> nameOFFirst, "age" -> ageOfFirst)
Map("name" -> nameOFsecond, "age" -> ageOfsecond)
Run Code Online (Sandbox Code Playgroud)

我尝试使用map via:x._1但这似乎不起作用Array [spark.sql.row]我怎么能进行转换呢?

scala apache-spark

0
推荐指数
1
解决办法
3691
查看次数