小编ssh*_*off的帖子

从Spark DataFrame中的单个列派生多个列

我有一个带有巨大可解析元数据的DF作为Dataframe中的单个字符串列,我们称之为DFA,使用ColmnA.

我想打破这一列,将ColmnA分成多个列,通过一个函数,ClassXYZ = Func1(ColmnA).此函数返回一个具有多个变量的类ClassXYZ,现在每个变量都必须映射到新列,例如ColmnA1,ColmnA2等.

如何通过调用此Func1一次,使用这些附加列从一个Dataframe到另一个Data转换,而不必重复它来创建所有列.

如果我每次都要调用这个巨大的函数添加一个新列,它很容易解决,但这是我希望避免的.

请使用工作或伪代码建议.

谢谢

桑杰

scala user-defined-functions dataframe apache-spark apache-spark-sql

48
推荐指数
3
解决办法
5万
查看次数

使用空/空字段值创建新的Dataframe

我正在从现有数据框架创建一个新的Dataframe,但需要在这个新DF中添加新列(下面代码中的"field1").我该怎么办?工作示例代码示例将不胜感激.

val edwDf = omniDataFrame 
  .withColumn("field1", callUDF((value: String) => None)) 
  .withColumn("field2",
    callUdf("devicetypeUDF", (omniDataFrame.col("some_field_in_old_df")))) 

edwDf
  .select("field1", "field2")
  .save("odsoutdatafldr", "com.databricks.spark.csv"); 
Run Code Online (Sandbox Code Playgroud)

scala dataframe apache-spark apache-spark-sql

25
推荐指数
2
解决办法
5万
查看次数

如何过滤Spark Dataframe的MapType字段

我有一个Spark Dataframe,其中一个字段是MapType ....我可以获取maptype字段的任何键的数据,但是当我为特定键的特定值应用过滤器时无法做到...

val line = List (("Sanjay", Map("one" -> 1, "two" -> 2)), ("Taru", Map("one" -> 10, "two" -> 20)) )
Run Code Online (Sandbox Code Playgroud)

我创建了上面列表的RDD和DF并尝试获取DF,Map值,其中值> = 5 .....但我在Spark Repl中得到以下异常..请帮助

val rowrddDFFinal = rowrddDF.select(rowrddDF("data.one").alias("data")).filter(rowrddDF("data.one").geq(5))
Run Code Online (Sandbox Code Playgroud)

org.apache.spark.sql.AnalysisException:已解析的属性数据#1 missin // | g来自运营商的数据#3!过滤器(数据#1 [one] AS one#4> = 5); // | 在org.apache.spark.sql.catalyst.analysis.CheckAnalysis $ class.failAnalys // | 是(CheckAnalysis.scala:38)// | 在org.apache.spark.sql.catalyst.analysis.Analyzer.failAnalysis(Analyzer // | .scala:42)// | 在org.apache.spark.sql.catalyst.analysis.CheckAnalysis $$ anonfun $ checkAn // | alysis $ 1.适用(CheckAnalysis.scala:121)// | 在org.apache.spark.sql.catalyst.analysis.CheckAnalysis $$ anonfun $ checkAn // | alysis $ 1.apply(CheckAnalysis.scala:50)// | 在org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala // |:98)// …

scala dataframe apache-spark apache-spark-sql

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

如何验证Spark Dataframe的内容

我有Scala Spark代码库,它运行良好,但不应该.

第二列有混合类型的数据,而在Schema我定义了它IntegerType.我的实际程序有超过100列,并DataFrames在转换后继续派生多个子项.

如何验证RDDDataFrame字段的内容是否具有正确的数据类型值,从而忽略无效行或将列的内容更改为某个默认值.任何更多的数据质量检查指针DataFrameRDD赞赏.

var theSeq = Seq(("X01", "41"),
    ("X01", 41),
    ("X01", 41),
    ("X02", "ab"),
    ("X02", "%%"))

val newRdd = sc.parallelize(theSeq)
val rowRdd = newRdd.map(r => Row(r._1, r._2))

val theSchema = StructType(Seq(StructField("ID", StringType, true),
    StructField("Age", IntegerType, true)))
val theNewDF = sqc.createDataFrame(rowRdd, theSchema)
theNewDF.show()  
Run Code Online (Sandbox Code Playgroud)

validation scala dataframe apache-spark apache-spark-sql

4
推荐指数
2
解决办法
9244
查看次数