我有以下代码:
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)
我怎么解决这个问题 ?
如何从 python 中获得 100000 个寄存器在 elasticsearch 中?MatchAll 查询仅检索 10000。
我使用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)
但我不想做以前的事情.
亲切的问候.
我正在处理一个文本文件,并将转换后的行从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
由于某些奇怪的原因,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)
谁能解释,为什么星火如此行事?
我在数据框中有 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 中的累积总和
我想转换一个创建的数组,如:
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]我怎么能进行转换呢?