【问题标题】:cant perform 2 succesive groupBy in spark无法在火花中执行 2 个连续的 groupBy
【发布时间】:2015-02-26 23:23:01
【问题描述】:

我正在 python 上使用 Spark。

我的问题是:我有一个 .csv 文件,其中包含一些数据(int1、int2、int3、日期)。我在int1 上做了一个groupByKey。现在我想用第一个groupBy创建的rdd在我的日期执行另一个groupBy

问题是我无法执行它。有什么想法吗?

问候

编辑2: 从 pyspark 导入 SparkContext 导入 csv 导入系统 导入字符串IO

sc = SparkContext("local", "Simple App")
file = sc.textFile("histories_2week9.csv")

 csvById12Rdd=file.map(lambda (id1,id2,value): ((id1,id2),value)).groupByKey()
 csvById1Rdd=csvById12Rdd.map(lambda ((id1,id2),group):(id1, (id2,group))).groupByKey()



def printit(one):
  id1, twos=one
  print("Id1:{}".format(id1))
    for two in twos:
      id2, values=two
      print("Id1:{} Id2:{}".format(id1,id2))
     for value in values:
        print("Id1:{} Id2:{} Value:{}".format(id1,id2,value))


  csvById12Rdd.first().foreach(printit)

csv 就像 31705,48,2,2014-10-28T18:14:09.000Z

编辑 3:

我可以用这段代码打印我的迭代器数据

from pyspark import SparkContext

import csv
import sys
import StringIO

sc = SparkContext("local", "Simple App")
file = sc.textFile("histories_2week9.csv")

def go_in_rdd2(x):
  print x[0]
  for i in x[1]:
      print i

counts = file.map(lambda line: (line.split(",")[0],line.split(",")[1:]))
counts = counts.groupByKey()
counts.foreach(go_in_rdd2)

但我仍然无法分组

【问题讨论】:

  • 是否有错误消息,您的工作是否崩溃,我们需要更多信息来回答问题

标签: python csv apache-spark


【解决方案1】:

Group by 返回一个 (Key, Iterable[Value]) 的 RDD,你能反过来做吗?

  1. 按 id1 id2 分组,得到 ((Id1,Id2), Iterable[Value]) 的 RDD
  2. 然后按id1 单独分组,得到(Id1,Iterable[(Id2,Iterable[Value])])的RDD

类似:

csv=[(1,1,"One","Un"),(1,2,"Two","Deux"),(2,1,"Three","Trois"),(2,1,"Four","Quatre")]
csvRdd=sc.parallelize(csv)
# Step 1
csvById12Rdd=csvRdd.map(lambda (id1,id2,value1,value2): ((id1,id2),(value1,value2))).groupByKey()
# Step 2
csvById1Rdd=csvById12Rdd.map(lambda ((id1,id2),group):(id1, (id2,group))).groupByKey()
# Print    
def printit(one):
    id1, twos=one
    print("Id1:{}".format(id1))
    for two in twos:
        id2, values=two
        print("Id1:{} Id2:{}".format(id1,id2))
        for value1,value2 in values:
            print("Id1:{} Id2:{} Values:{} {}".format(id1,id2,value1,value2))

csvById1Rdd.foreach(printit)

【讨论】:

  • 感谢您的宝贵时间。当我尝试你的方法时,我在调用你的打印函数时收到一条错误消息,我得到“ValueError:解包的值太多”但我认为 2 groupBy 有效
  • 我已经在 Spark 1.2.1 和 Python 2.7 中进行了测试。你使用 Python 3 吗?你接受这个答案吗?
  • 我也使用 Python 2.7 和 spark 1.2.1
  • 我猜你在你的元组中添加了一些值。 “ValueError: too many values to unpack”表示元组的大小不合适,例如,如果您写了类似“one, two=(1,2,3)”的内容(左侧有两个元素,左侧有三个右边三个)
  • 我正在尝试修改您的代码以使其适用于具有 4 个字段的 csv,您认为这是个好主意吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-11-24
  • 2020-07-09
  • 2022-06-11
  • 1970-01-01
  • 1970-01-01
  • 2019-05-02
  • 1970-01-01
相关资源
最近更新 更多