【问题标题】:dataframe from hive table to iterate through each element for some operation and write in df,rdd,list来自 hive 表的数据帧以迭代每个元素以进行某些操作并写入 df,rdd,list
【发布时间】:2019-04-17 23:23:54
【问题描述】:

我有一个DF,输入数据如下:

+----+----+
|col1|col2|
+----+--------+
| abc|2E2J2K2F|
| bcd|    2K3D|
+----+--------+

我预期的预期输出是:

+-----+-----+
| col1| col2|
+----+------+
| abc|    2E|
| abc|    2J|
| abc|    2K|
| abc|    2F|
| bcd|    2K|
| bcd|    3D|
+----+------+
+----+------+

【问题讨论】:

  • |col1|col2| +----+--------+ | abc|2E2J2K2F| |光盘| 2K3D| +----+--------+ 预期输出 +-----+-----+ | col1| col2| +----+------+ | ABC| 2E| | ABC| 2J| | ABC| 2K| | ABC| 2F| |光盘| 2K| |光盘| 3D| +----+------+ +----+------+

标签: scala list apache-spark dataframe rdd


【解决方案1】:

使用 udf() 分割字符串,然后将其分解。看看这个:

scala>  val df = Seq(("abc","2E2J2K2F"),("bcd","2K3D")).toDF("col1","col2")
df: org.apache.spark.sql.DataFrame = [col1: string, col2: string]

scala> def split2(x:String):Array[String] = x.sliding(2,2).toArray
split2: (x: String)Array[String]

scala> val myudf_split2 = udf ( split2(_:String):Array[String] )
myudf_split2: org.apache.spark.sql.expressions.UserDefinedFunction = UserDefinedFunction(<function1>,ArrayType(StringType,true),Some(List(StringType)))

scala> df.withColumn("newcol",explode(myudf_split2('col2))).select("col1","newcol").show
+----+------+
|col1|newcol|
+----+------+
| abc|    2E|
| abc|    2J|
| abc|    2K|
| abc|    2F|
| bcd|    2K|
| bcd|    3D|
+----+------+


scala>

更新:

split2() 只是将字符串拆分为 2 个字节并创建一个数组。 分解函数根据数组的长度复制行,为所有行提供每个索引值。

scala> def split2(x:String):Array[String] = x.sliding(2,2).toArray
split2: (x: String)Array[String]

scala> split2("12345678")
res168: Array[String] = Array(12, 34, 56, 78)

scala> def split2(x:String):Array[String] = x.sliding(2,2).toArray
split2: (x: String)Array[String]

scala> split2("12345678")
res168: Array[String] = Array(12, 34, 56, 78)

scala> "12345678".sliding(4,4).toArray
res171: Array[String] = Array(1234, 5678)

【讨论】:

  • 错误:类型不匹配; [错误] 发现:需要符号 [错误]:org.apache.spark.sql.Column [错误] df.withColumn("newcol",explode(myudf_split2('col2))).select("col1","newcol" ).show [ERROR] ^ [ERROR] 发现一个错误我有 scala 2.12.x 和 spark 1.6.x(工作环境)
  • 你的 df.schema 是什么?
  • spark 1.6 old..我不认为它支持 udf 函数..请考虑使用 spark 2.x 版本
  • 抱歉,因为我们有 CDH 5.14,所以不能这样做。包裹,如果有任何替代方案,请提出建议
  • 这个 cloudera 版本可能有 spark 2.x.. 你可能正在启动“spark-shell”,请尝试使用“spark2-shell”启动。另外,您在哪个步骤中遇到错误?是 udf() 吗?
【解决方案2】:

val df = Seq(("abc","2E2J2K2F"),("bcd","2K3D")).toDF("col1","col2") df: org.apache.spark.sql.DataFrame = [col1: string, col2: string]

scala> def split2(x:String):Array[String] = x.sliding(2,2).toArray split2: (x: String)Array[String]

scala> val myudf_split2 = udf ( split2(_:String):Array[String] ) myudf_split2: org.apache.spark.sql.expressions.UserDefinedFunction = UserDefinedFunction(,ArrayType(StringType,true),Some(List(StringType))))

scala> df.withColumn("newcol",explode(myudf_split2(df.col("col2")))).select("col1","newcol").show

+----+------+ |col1|新科尔| +----+------+ | ABC| 2E| | ABC| 2J| | ABC| 2K| | ABC| 2F| |光盘| 2K| |光盘| 3D| +----+-----+

【讨论】:

    猜你喜欢
    • 2020-09-07
    • 2020-01-05
    • 1970-01-01
    • 2020-05-09
    • 1970-01-01
    • 2018-11-09
    • 2014-06-18
    • 2016-09-27
    • 2020-12-12
    相关资源
    最近更新 更多