val columnName=Seq("col1","col2",....."coln");
Run Code Online (Sandbox Code Playgroud)
有没有办法执行dataframe.select操作以获取仅包含指定列名的数据帧.我知道我可以做,dataframe.select("col1","col2"...)
但是columnName在运行时生成.我可以dataframe.select()在循环中为每个列名重复执行.它会有任何性能开销吗?有没有其他更简单的方法来实现这一目标?
我有一个具有以下结构的数据帧:
|-- data: struct (nullable = true)
| |-- id: long (nullable = true)
| |-- keyNote: struct (nullable = true)
| | |-- key: string (nullable = true)
| | |-- note: string (nullable = true)
| |-- details: map (nullable = true)
| | |-- key: string
| | |-- value: string (valueContainsNull = true)
Run Code Online (Sandbox Code Playgroud)
如何展平结构并创建新的数据框:
|-- id: long (nullable = true)
|-- keyNote: struct (nullable = true)
| |-- key: string (nullable = true)
| |-- note: …Run Code Online (Sandbox Code Playgroud) 我有一个DataFrame与Timestamp列,我需要为转换Date格式.
是否有可用的Spark SQL函数?
我想在DataSet中为Row类型编写一个编码器,用于我正在进行的地图操作.基本上,我不明白如何编写编码器.
以下是地图操作的示例:
In the example below, instead of returning Dataset<String>, I would like to return Dataset<Row>
Dataset<String> output = dataset1.flatMap(new FlatMapFunction<Row, String>() {
@Override
public Iterator<String> call(Row row) throws Exception {
ArrayList<String> obj = //some map operation
return obj.iterator();
}
},Encoders.STRING());
Run Code Online (Sandbox Code Playgroud)
我明白,编码器需要编写如下代码:
Encoder<Row> encoder = new Encoder<Row>() {
@Override
public StructType schema() {
return join.schema();
//return null;
}
@Override
public ClassTag<Row> clsTag() {
return null;
}
};
Run Code Online (Sandbox Code Playgroud)
但是,我不理解编码器中的clsTag(),我试图找到一个可以演示相似内容的运行示例(即行类型的编码器)
编辑 - 这不是所提问题的副本:尝试将数据帧行映射到更新行时编码器错误,因为答案谈到在Spark 2.x中使用Spark 1.x(我不是这样做),我也在寻找用于Row类的编码器而不是解决错误.最后,我一直在寻找Java解决方案,而不是Scala.
java apache-spark apache-spark-sql apache-spark-dataset apache-spark-encoders
我通过Spark 1.5.0使用PySpark.对于datetime值,我在列的行中有一个不常见的String格式.它看起来像这样:
Row[(daytetime='2016_08_21 11_31_08')]
Run Code Online (Sandbox Code Playgroud)
有没有办法将这种非正统yyyy_mm_dd hh_mm_dd格式转换为时间戳?最终可能出现的问题
df = df.withColumn("date_time",df.daytetime.astype('Timestamp'))
Run Code Online (Sandbox Code Playgroud)
我原以为像星火SQL函数regexp_replace可以工作,但我当然需要更换
_与-在日期一半_用:在部分时间.
我想我可以在2中拆分列,substring并从时间结束后向后计数.然后单独执行'regexp_replace',然后连接.但这似乎很多操作?有没有更简单的方法?
我有DataFrame一些列.现在我想在现有的DataFrame中再添加两列.
目前我正在使用withColumnDataFrame中的方法.
例如:
df.withColumn("newColumn1", udf(col("somecolumn")))
.withColumn("newColumn2", udf(col("somecolumn")))
Run Code Online (Sandbox Code Playgroud)
实际上我可以使用Array [String]在单个UDF方法中返回两个newcoOlumn值.但目前这就是我的做法.
无论如何,我能有效地做到这一点吗?使用explode是不错的选择?
即使我必须使用explode,我必须使用withColumn一次,然后返回列值Array[String],然后使用explode,再创建两列.
哪一个有效?还是有其他选择吗?
我不知道为什么我会遇到困难,看起来很简单,因为在R或熊猫中相当容易.我想避免使用pandas,因为我正在处理大量数据,我相信toPandas()所有数据都会加载到pyspark中的驱动程序内存中.
我有2个数据帧:df1和df2.我想过滤df1(删除所有行)df1.userid = df2.useridAND df1.group = df2.group.我不知道我是否应该使用filter(),join()或sql 例如:
df1:
+------+----------+--------------------+
|userid| group | all_picks |
+------+----------+--------------------+
| 348| 2|[225, 2235, 2225] |
| 567| 1|[1110, 1150] |
| 595| 1|[1150, 1150, 1150] |
| 580| 2|[2240, 2225] |
| 448| 1|[1130] |
+------+----------+--------------------+
df2:
+------+----------+---------+
|userid| group | pick |
+------+----------+---------+
| 348| 2| 2270|
| 595| 1| 2125|
+------+----------+---------+
Result I want:
+------+----------+--------------------+ …Run Code Online (Sandbox Code Playgroud) 在Spark 2.1 文档中提到了这一点
Spark运行在Java 7 +,Python 2.6 +/3.4 +和R 3.1+上.对于Scala API,Spark 2.1.0使用Scala 2.11.您需要使用兼容的Scala版本(2.11.x).
在Scala 2.12 发布消息时,它还提到:
尽管Scala 2.11和2.12主要是源兼容的,以便于交叉构建,但它们不是二进制兼容的.这使我们能够不断改进Scala编译器和标准库.
但是当我构建一个超级jar(使用Scala 2.12)并在Spark 2.1上运行它时.一切都很好.
我知道它不是任何官方消息来源,但在47度博客中他们提到Spark 2.1确实支持Scala 2.12.
如何解释那些(冲突?)信息?
我需要使用Spark SQL从Hive表加载数据HiveContext并加载到HDFS中.默认情况下,DataFramefrom SQL输出有2个分区.为了获得更多的并行性,我需要更多的SQL分区.HiveContext中没有重载方法来获取分区数参数.
RDD的重新分区导致改组并导致更多的处理时间.
>
val result = sqlContext.sql("select * from bt_st_ent")
Run Code Online (Sandbox Code Playgroud)
有日志输出:
Starting task 0.0 in stage 131.0 (TID 297, aster1.com, partition 0,NODE_LOCAL, 2203 bytes)
Starting task 1.0 in stage 131.0 (TID 298, aster1.com, partition 1,NODE_LOCAL, 2204 bytes)
Run Code Online (Sandbox Code Playgroud)
我想知道有没有办法增加SQL输出的分区大小.
Spark 2.2引入了Kafka的结构化流媒体源.据我所知,它依靠HDFS检查点目录来存储偏移并保证"完全一次"的消息传递.
但旧的码头(如https://blog.cloudera.com/blog/2017/06/offset-management-for-apache-kafka-with-apache-spark-streaming/)表示Spark Streaming检查点无法跨应用程序恢复或Spark升级,因此不太可靠.作为一种解决方案,有一种做法是支持在支持MySQL或RedshiftDB等事务的外部存储中存储偏移量.
如果我想将Kafka源的偏移存储到事务DB,我如何从结构化流批处理中获得偏移量?
以前,可以通过将RDD转换为HasOffsetRanges:
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
Run Code Online (Sandbox Code Playgroud)
但是使用新的Streaming API,我有一个Dataset,InternalRow我找不到一个简单的方法来获取偏移量.Sink API只有addBatch(batchId: Long, data: DataFrame)方法,我怎么能想得到给定批次ID的偏移量?
offset apache-kafka apache-spark apache-spark-sql spark-structured-streaming
apache-spark ×9
dataframe ×4
scala ×3
java ×2
pyspark ×2
abi ×1
apache-kafka ×1
hive ×1
offset ×1
partitioning ×1
python-2.7 ×1
timestamp ×1