【问题标题】:How do I split a column by using delimiters from another column in Spark/Scala如何使用 Spark/Scala 中另一列的分隔符拆分列
【发布时间】:2021-09-23 14:45:19
【问题描述】:

我还有一个与拆分功能有关的问题。 我是 Spark/Scala 的新手。

下面是示例数据框 -


+-------------------+---------+
|             VALUES|Delimiter|
+-------------------+---------+
|       50000.0#0#0#|        #|
|          0@1000.0@|        @|
|                 1$|        $|
|1000.00^Test_string|        ^|
+-------------------+---------+

我希望输出为 -

+-------------------+---------+----------------------+
|VALUES             |Delimiter|split_values          |
+-------------------+---------+----------------------+
|50000.0#0#0#       |#        |[50000.0, 0, 0, ]     |
|0@1000.0@          |@        |[0, 1000.0, ]         |
|1$                 |$        |[1, ]                 |
|1000.00^Test_string|^        |[1000.00, Test_string]|
+-------------------+---------+----------------------+

我尝试手动拆分 -

dept.select(split(col("VALUES"),"#|@|\\$|\\^").show()

输出是 -

+-----------------------+
|split(VALUES,#|@|\$|\^)|
+-----------------------+
|      [50000.0, 0, 0, ]|
|          [0, 1000.0, ]|
|                  [1, ]|
|   [1000.00, Test_st...|
+-----------------------+


但我想为大型数据集自动拉出分隔符。

【问题讨论】:

  • 请在问题中添加您尝试过的内容和失败的内容。您应该为您的问题提供一个可重复的最小示例。
  • 我知道可以根据上面的示例数据框手动完成 - ``` dept.select(split(col("VALUES"),"#|@|\\$|\ \^").show() ``` 和输出匹配,但我不想手动放置分隔符。
  • 额外列的逻辑是什么?
  • @dsk 我已经编辑了这个问题。到目前为止,我不确定额外列的逻辑。我主要关心的是自动获取大型数据集的分隔符。
  • @Shashwat - 请参阅下面的答案 - 如果您对答案感到满意,请不要犹豫接受并投票 :)

标签: scala apache-spark apache-spark-sql


【解决方案1】:

您需要使用 expr 和 split() 来进行动态拆分

df = spark.createDataFrame([("50000.0#0#0#","#"),("0@1000.0@","@")],["VALUES","Delimiter"])
df = df.withColumn("split", F.expr("""split(VALUES, Delimiter)"""))
df.show()

+------------+---------+-----------------+
|      VALUES|Delimiter|            split|
+------------+---------+-----------------+
|50000.0#0#0#|        #|[50000.0, 0, 0, ]|
|   0@1000.0@|        @|    [0, 1000.0, ]|
+------------+---------+-----------------+

【讨论】:

  • 这是一个非常酷的答案。 F 的导入是什么?是regex.F吗?我无法让它工作。
  • 嗨 @PubuduSitinamaluwa 函数 expr() 在 spark.sql.functions 下。您可以直接导入 expr 或使用别名调用 expr()。希望这会有所帮助。
  • 嗨,@dsk 我认为我们必须在其他情况下包含转义字符才能处理“$”、“^”等。
  • @PubuduSitinamaluwa from pyspark.sql import functions as F from pyspark.sql import types as T from pyspark.sql import Window as W
  • @Shashwat 如果您能接受答案,将不胜感激 :) 很高兴它对您有所帮助
【解决方案2】:

编辑:请检查 scala 版本的答案底部。

您可以使用自定义的用户定义函数 (pyspark.sql.functions.udf) 来实现此目的。

from typing import List

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType, ArrayType


def split_col(value: StringType, delimiter: StringType) -> List[str]:
    return str(value).split(str(delimiter))


udf_split = udf(lambda x, y: split_col(x, y), ArrayType(StringType()))

spark = SparkSession.builder.getOrCreate()

df = spark.createDataFrame([
    ('50000.0#0#0#', '#'), ('0@1000.0@', '@'), ('1$', '$'), ('1000.00^Test_string', '^')
], schema='VALUES String, Delimiter String')

df = df.withColumn("split_values", udf_split(df['VALUES'], df['Delimiter']))

df.show(truncate=False)

输出

+-------------------+---------+----------------------+
|VALUES             |Delimiter|split_values          |
+-------------------+---------+----------------------+
|50000.0#0#0#       |#        |[50000.0, 0, 0, ]     |
|0@1000.0@          |@        |[0, 1000.0, ]         |
|1$                 |$        |[1, ]                 |
|1000.00^Test_string|^        |[1000.00, Test_string]|
+-------------------+---------+----------------------+

请注意,split_values 列包含一个字符串列表。您还可以更新split_col 函数以对值进行更多更改。

编辑: Scala 版本

import org.apache.spark.sql.functions.udf

import spark.implicits._

val data = Seq(("50000.0#0#0#", "#"), ("0@1000.0@", "@"), ("1$", "$"), ("1000.00^Test_string", "^"))
var df = data.toDF("VALUES", "Delimiter")

val udf_split_col = udf {(x:String,y:String)=> x.split(y)}

df = df.withColumn("split_values", udf_split_col(df.col("VALUES"), df.col("Delimiter")))

df.show(false)

编辑 2

为避免正则表达式中使用特殊字符的问题,您可以在使用split() 方法时使用char 而不是String,如下所示。

val udf_split_col = udf { (x: String, y: String) => x.split(y.charAt(0)) }

【讨论】:

  • OP 要求使用 Scala ;)
  • 感谢您指出这一点。但概念是一样的。
  • 当然,但不是每个人都能将 Python“翻译”成 Scala。
  • 我同意。让我想出一个 scala 版本。
  • @Shashwat 希望您能够解决问题。我添加了一个 scala 版本的示例。
【解决方案3】:

这是另一种处理方式,使用 sparksql

df.createOrReplaceTempView("test")

spark.sql("""select VALUES,delimiter,split(values,case when delimiter in ("$","^") then concat("\\",delimiter) else delimiter end) as split_value from test""").show(false)

请注意,我包含了 case when 语句来添加转义字符来处理“$”和“^”的大小写,否则它不会拆分。

+-------------------+---------+----------------------+
|VALUES             |delimiter|split_value           |
+-------------------+---------+----------------------+
|50000.0#0#0#       |#        |[50000.0, 0, 0, ]     |
|0@1000.0@          |@        |[0, 1000.0, ]         |
|1$                 |$        |[1, ]                 |
|1000.00^Test_string|^        |[1000.00, Test_string]|
+-------------------+---------+----------------------+

【讨论】:

    【解决方案4】:

    这是我最近的解决方案

    import java.util.regex.Pattern
    val split_udf = udf((value: String, delimiter: String) => value.split(Pattern.quote(delimiter), -1))
    val solution = dept.withColumn("split_values", split_udf(col("VALUES"),col("Delimiter")))
    solution.show(truncate = false)
    

    它将跳过分隔符列中的特殊字符。 其他答案不适用于

    ("50000.0\\0\\0\\", "\\")
    

    而 linusRian 的回答需要手动添加特殊字符

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-11-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多