【问题标题】:How do I remove words numerics pyspark如何删除单词数字pyspark
【发布时间】:2019-05-25 20:20:04
【问题描述】:

我试图从我的单词数组中仅删除数字单词,但我创建的函数无法正常工作。当我尝试从我的数据框中查看信息时,会出现以下错误消息

首先我转换了我的字符串和单词标记

from pyspark.ml.feature import RegexTokenizer
regexTokenizer = RegexTokenizer(
    inputCol="description",
    outputCol="words_withnumber",
    pattern="\\W"
)

data = regexTokenizer.transform(data)

我创建了只删除数字的函数

from pyspark.sql.functions import when,udf
from pyspark.sql.types import BooleanType

def is_digit(value):
    if value:
        return value.isdigit()
    else:
        return False

is_digit_udf = udf(is_digit, BooleanType())

调用函数

data = data.withColumn(
    'words_withoutnumber', 
    when(~is_digit_udf(data['words_withnumber']), data['words_withnumber'])
)

错误:

org.apache.spark.SparkException:作业因阶段失败而中止:阶段 5.0 中的任务 0 失败 4 次,最近一次失败:阶段 5.0 中丢失任务 0.3(TID 14、10.139.64.4、执行程序 0):org .apache.spark.api.python.PythonException:回溯(最近一次调用最后一次):

示例数据框

+-----------+-----------------------------------------------------------+
|categoryid |description                                                |
+-----------+-----------------------------------------------------------+
|      33004|["short","sarja", "40567","detalhe","couro"]               | 
|      22033|["multipane","6768686868686867868888","220v","branco"]     | 
+-----------+-----------------------------------------------------------+

预期结果

+-----------+-----------------------------------------------------------+
|categoryid |description                                                |
+-----------+-----------------------------------------------------------+
|      33004|["short","sarja","detalhe","couro"]                        | 
|      22033|["multipane","220v","branco"]                              |
+-----------+-----------------------------------------------------------+

【问题讨论】:

  • 您 udf 需要一个字符串,但您将一个数组传递给它。同样在您的示例数据框中,description 应该是 words_withnumber?
  • 您需要iterate over the array 并过滤掉所需的单词。你用的是什么版本的火花?
  • @Psidom,我尝试循环遍历数组,但收到以下错误消息“未定义名称'ArrayType'
  • @pault 版本火花 2.3.1。我看到链接,我试试这个 filter_udf = udf(lambda row: [x for x in row if is_digit(x)], ArrayType(StringType())) ,我收到错误“名称'ArrayType'未定义” – user3661384 25分钟前
  • from pyspark.sql.types import ArrayType

标签: pyspark apache-spark-sql


【解决方案1】:

既然您提到了 pyspark,那么让我们包含 python。

数据帧

data =[ (33004     , ['short', 'sarja', '40567', 'detalhe', 'couro']),
       (22033     , ['multipane', '6768686868686867868888', '220v', 'branco']),
       (33033    ,['123', 'x220', '220v', '889'] )  
]

df= spark.createDataFrame(data, ('categoryid','description' ))

代码

#仅限字母数字

df.withColumn("description_filtered", expr("filter(description, x -> x rlike '([a-z]+)')")).show()

结果

+----------+--------------------+--------------------+
|categoryid|         description|description_filtered|
+----------+--------------------+--------------------+
|     33004|[short, sarja, 40...|[short, sarja, de...|
|     22033|[multipane, 67686...|[multipane, 220v,...|
|     33033|[123, x220, 220v,...|        [x220, 220v]|
+----------+--------------------+--------------------+

#仅按字母顺序排列

df.withColumn("description_filtered", expr("filter(description, x -> x rlike '([^0-9]{2})')")).show()

结果

+----------+--------------------+--------------------+
|categoryid|         description|description_filtered|
+----------+--------------------+--------------------+
|     33004|[short, sarja, 40...|[short, sarja, de...|
|     22033|[multipane, 67686...| [multipane, branco]|
|     33033|[123, x220, 220v,...|                  []|
+----------+--------------------+--------------------+

【讨论】:

    【解决方案2】:

    如果您出于性能原因想避免使用 udf() 并且逗号不会出现在您的“描述”列中,那么下面的 scala 解决方案将起作用。 df.withColumn() 在 pyspark 中应该是类似的。

    注意:我还添加了第三条记录,以表明当数字出现在数组的开头/结尾时,该解决方案有效。试试看。

    scala> val df = Seq((33004,Array("short","sarja", "40567","detalhe","couro")), (22033,Array("multipane","6768686868686867868888","220v","branco")), (33033,Array("0123","x220","220v","889"))).toDF("categoryid","description")
    df: org.apache.spark.sql.DataFrame = [categoryid: int, description: array<string>]
    
    scala> df.show(false)
    +----------+-------------------------------------------------+
    |categoryid|description                                      |
    +----------+-------------------------------------------------+
    |33004     |[short, sarja, 40567, detalhe, couro]            |
    |22033     |[multipane, 6768686868686867868888, 220v, branco]|
    |33033     |[0123, x220, 220v, 889]                          |
    +----------+-------------------------------------------------+
    
    
    scala> df.withColumn("newc",split(regexp_replace(regexp_replace(regexp_replace(concat_ws(",",'description),"""\b\d+\b""",""),"""^,|,$""",""),",,",","),",")).show(false)
    +----------+-------------------------------------------------+------------------------------+
    |categoryid|description                                      |newc                          |
    +----------+-------------------------------------------------+------------------------------+
    |33004     |[short, sarja, 40567, detalhe, couro]            |[short, sarja, detalhe, couro]|
    |22033     |[multipane, 6768686868686867868888, 220v, branco]|[multipane, 220v, branco]     |
    |33033     |[0123, x220, 220v, 889]                          |[x220, 220v]                  |
    +----------+-------------------------------------------------+------------------------------+
    
    
    scala>
    

    Spark 2.4 答案

    在2.4版本以后使用spark-sql,可以使用filter()高阶函数并得到结果

    scala> val df = Seq((33004,Array("short","sarja", "40567","detalhe","couro")), (22033,Array("multipane","6768686868686867868888","220v","branco")), (33033,Array("0123","x220","220v","889"))).toDF("categoryid","description")
    df: org.apache.spark.sql.DataFrame = [categoryid: int, description: array<string>]
    
    scala> df.createOrReplaceTempView("tab")
    
    scala> spark.sql(""" select categoryid, filter(description, x -> lower(x)!=upper(x)) fw from tab """).show(false)
    +----------+------------------------------+
    |categoryid|fw                            |
    +----------+------------------------------+
    |33004     |[short, sarja, detalhe, couro]|
    |22033     |[multipane, 220v, branco]     |
    |33033     |[x220, 220v]                  |
    +----------+------------------------------+
    
    
    scala>
    

    【讨论】:

      【解决方案3】:

      作为@pault 的帮助,解决方案是这样的。

      from pyspark.sql.functions import when,udf
      from pyspark.sql.types import BooleanType
      
      def is_digit(value):
          if value:
              return value.isdigit()
          else:
              return False
      
      is_digit_udf = udf(is_digit, BooleanType()
      

      调用函数

      from pyspark.sql.types import ArrayType, StringType
      from pyspark.sql.types import StructType
      
      filter_length_udf = udf(lambda row: [x for x in row if not is_digit(x)], ArrayType(StringType()))
      
      data = data.withColumn('words_clean', filter_length_udf(col('words_withnumber')))
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2018-04-10
        • 2023-03-22
        • 1970-01-01
        • 1970-01-01
        • 2016-04-08
        • 1970-01-01
        相关资源
        最近更新 更多