【问题标题】:How to filter pyspark dataframes如何过滤 pyspark 数据帧
【发布时间】:2021-01-08 04:08:23
【问题描述】:

我看到了许多与过滤 pyspark 数据帧相关的问题,但尽管我尽了最大努力,我还是无法让任何非 SQL 解决方案发挥作用。

+----------+-------------+-------+--------------------+--------------+---+
|purch_date|  purch_class|tot_amt|       serv-provider|purch_location| id|
+----------+-------------+-------+--------------------+--------------+---+
|03/11/2017|Uncategorized| -17.53|             HOVER  |              |  0|
|02/11/2017|    Groceries| -70.05|1774 MAC'S CONVEN...|     BRAMPTON |  1|
|31/10/2017|Gasoline/Fuel|    -20|              ESSO  |              |  2|
|31/10/2017|       Travel|     -9|TORONTO PARKING A...|      TORONTO |  3|
|30/10/2017|    Groceries|  -1.84|         LONGO'S # 2|              |  4|

这不起作用:

df1 = spark.read.csv("/some/path/to/file", sep=',')\
            .filter((col('purch_location')=='BRAMPTON')

这没有用

df1 = spark.read.csv("/some/path/to/file", sep=',')\
            .filter(purch_location == 'BRAMPTON')

这个(SQL 表达式)有效,但需要很长时间,我想有一个更快的非 SQL 方法

df1 = spark.read.csv("/some/path/to/file", sep=',')\
            .filter(purch_location == 'BRAMPTON')

更新我应该提到我能够使用类似的方法(比 SQL 表达式运行得更快):

df1 = spark.read.csv("/some/path/to/file", sep=',')
df2 = df1.filter(df1.purch_location == "BRAMPTON")

但想了解为什么“管道”/ 连接语法不正确。

【问题讨论】:

    标签: python apache-spark filter pyspark apache-spark-sql


    【解决方案1】:

    你可以使用df["purch_location"]

    df = spark.read.csv("/some/path/to/file", sep=',')
    df = df.filter(df["purch_location"] == "BRAMPTON")
    

    【讨论】:

    • 谢谢,我希望使用/ 连接器,我应该提到我能够在连接之外做到这一点
    【解决方案2】:

    如果你坚持使用反斜杠,你可以这样做:

    from pyspark.sql.functions import col
    
    df = spark.read.csv('/some/path/to/file', sep=',') \
         .filter(col('purch_location') == 'BRAMPTON')
    

    您的第一次尝试失败了,因为括号不平衡。

    另外,BRAMPTON 字符串后面似乎还有一些空格,所以你可能想先trim 列:

    from pyspark.sql.functions import col, trim
    
    df = spark.read.csv('/some/path/to/file', sep=',') \
         .filter(trim(col('purch_location')) == 'BRAMPTON')
    

    【讨论】:

    • 我收到以下错误:NameError: name 'col' is not defined
    • @CyrusMohammadian 你需要导入它,看看我编辑的答案?
    • 是的,我的错,我想我没有运行前一个单元格
    猜你喜欢
    • 2017-03-16
    • 2021-04-10
    • 2017-10-21
    • 1970-01-01
    • 2019-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-24
    相关资源
    最近更新 更多