【问题标题】:PySpark groupBy indexerror: list out of rangePySpark groupBy indexerror:列表超出范围
【发布时间】:2018-03-22 05:26:02
【问题描述】:

我的本​​地系统中有一个 csv 文件(包括标题),我试图在该文件上执行 groupBy 功能,即按目的分组和每个目的的金额总和。我在 pyspark 控制台上键入的命令如下

from pyspark import SparkContext, SparkConf
from pyspark.sql.types import *
from pyspark.sql import Row
csv_data=sc.textFile("/project/sample.csv").map(lambda p: p.split(",")) 
header = csv_data.first()
csv_data = csv_data.filter(lambda p:p != header)
df_csv  =  csv_data.map(lambda p: Row(checkin_acc  =  p[0],  duration = 
int(p[1]), credit_history = p[2], purpose = p[3], amount = int(p[4]),
svaing_acc = p[5], present_emp_since = p[6], inst_rate = int(p[7]), 
personal_status = p[8], other_debtors = p[9],residing_since = int(p[10]), 
property = p[11], age = int(p[12]), inst_plans = p[13], housing = p[14], 
num_credits = int(p[15]), job = p[16], dependents = int(p[17]), telephone = 
p[18], foreign_worker = p[19], status = p[20])).toDF()

grouped = df_csv.groupBy('purpose').sum('amount')
grouped.show()
[Stage 9:>                                                          (0 + 2) / 2]18/03/22 10:34:52 ERROR executor.Executor: Exception in task 1.0 in stage 9.0 (TID 10)
org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/opt/spark-2.2.1-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/worker.py", line 177, in main
    process()
  File "/opt/spark-2.2.1-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/worker.py", line 172, in process
    serializer.dump_stream(func(split_index, iterator), outfile)
  File "/opt/spark-2.2.1-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/serializers.py", line 268, in dump_stream
    vs = list(itertools.islice(iterator, batch))
  File "<stdin>", line 1, in <lambda>
IndexError: list index out of range

    at org.apache.spark.api.python.PythonRunner$$anon$1.read(PythonRDD.scala:193)
    at org.apache.spark.api.python.PythonRunner$$anon$1.next(PythonRDD.scala:156)
    at org.apache.spark.api.python.PythonRunner$$anon$1.next(PythonRDD.scala:152)
    at org.apache.spark.InterruptibleIterator.next(InterruptibleIterator.scala:40)
    at scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:4

如何解决这个错误?

【问题讨论】:

  • 试试grouped = df.groupby(['purpose','housing']).sum('amount')
  • 此错误表示您的 csv 文件中的字段少于 20 个。另外,我看到你想过滤掉标题行,所以试试这个: csv_data = spark.read.option('header','false').option('inferSchema','true').csv('/project /sample.csv') 此选项仅在您使用 spark 版本 2+ 时可用

标签: python python-3.x apache-spark pyspark apache-spark-sql


【解决方案1】:

如果你使用 pyspark 2+,你可以使用spark.read.csv

df = spark.read.csv ("project/sample.csv", header=True)

如果您想自己设置列名和类型,您也可以使用 StructType 定义架构并使用 schema kwarg 传递它。

【讨论】:

    【解决方案2】:

    IndexError: 列表索引超出范围

    上面的错误仅仅意味着当你使用分割文本文件的行时没有生成20个字段

    过滤解决方案

    一种方法是过滤掉所有长度小于20的行

    csv_data=sc.textFile("/project/sample.csv").map(lambda p: p.split(","))
    header = csv_data.first()
    csv_data = csv_data.filter(lambda p:p != header)\
        .filter(lambda x: len(x) == len(header))    #filter added 
    df_csv  =  csv_data.map(lambda p: Row(checkin_acc=p[0],
                                          duration=int(p[1]),
                                          credit_history=p[2],
                                          purpose=p[3],
                                          amount=int(p[4]),
                                          ..... #the rest of the codes are same
    

    添加虚拟数据

    另一种解决方案是在长度小于20的行中添加虚拟数据

    #function definition for adding dummy strings in case of less fields in data
    def addDummy(arr, header):
        headerLength = len(header)
        arrayLength = len(arr)
        if arrayLength > headerLength:
            return arr[:headerLength-1]
        elif arrayLength < headerLength:
            return arr + ["dummy" for x in range(0, headerLength-arrayLength)]
        else:
            return arr
    
    csv_data=sc.textFile("/project/sample.csv").map(lambda p: p.split(","))
    header = csv_data.first()
    csv_data = csv_data.filter(lambda p:p != header)\
        .map(lambda p: addDummy(p, header))   #map function added for checking length and adding dummy string in case of less fields
    df_csv  =  csv_data.map(lambda p: Row(checkin_acc=p[0],
                                          duration=int(p[1]),
                                          credit_history=p[2],
                                          purpose=p[3],
                                          amount=int(p[4]),
                                          svaing_acc=p[5],
                                          .... #the rest of the codes are same as in the question
    

    【讨论】:

      猜你喜欢
      • 2018-10-31
      • 2021-11-27
      • 1970-01-01
      • 2019-09-08
      • 2011-10-31
      • 2015-06-26
      • 1970-01-01
      相关资源
      最近更新 更多