使用 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)