【问题标题】:Python Spark How to Map Fields of one rdd to another rddPython Spark如何将一个rdd的字段映射到另一个rdd
【发布时间】:2015-12-10 09:46:19
【问题描述】:

根据上述主题,我对 python spark 非常陌生,我想将一个 Rdd 的字段映射到另一个 Rdd 的字段。示例如下

rdd1:

c_id    name 
121210  abc
121211  pqr

rdd2:

c_id   cn_id cn_value
121211  0     0
121210  0     1

所以匹配的 c_id 将被 name 替换为 cnid 和聚合的 cn_value。所以输出会像这样 abc 0 0 pqr 0 1

from pyspark import SparkContext
import csv
sc = SparkContext("local", "spark-App")
file1 = sc.textFile('/home/hduser/sample.csv').map(lambda line:line.split(',')).filter(lambda line:len(line)>1)
file2 = sc.textFile('hdfs://localhost:9000/sample2/part-00000').map(lambda line:line.split(','))
file1_fields = file1.map(lambda x: (x[0],x[1]))
file2_fields = file2.map(lambda x: (x[0],x[1],float(x[2])))

我怎样才能通过在这里放一些代码来实现我的目标。

任何帮助将不胜感激 谢谢你

【问题讨论】:

  • 有一个叫做join的操作和一个叫做DataFrame的结构。你也应该看看spark-csv
  • @mlk 我想给 OP 一个展示努力的机会 :)

标签: python apache-spark pyspark


【解决方案1】:

您要查找的操作称为join。鉴于您的结构,最好使用DataFrames 和spark-csv(我假设第二个文件也是逗号分隔的,但没有标题)。让我们从虚拟数据开始:

file1 = ... # path to the first file
file2 = ... # path to the second file

with open(file1, "w") as fw:
    fw.write("c_id,name\n121210,abc\n121211,pqr")

with open(file2, "w") as fw:
    fw.write("121211,0,0\n121210,0,1")

读取第一个文件:

df1 = (sqlContext.read 
    .format('com.databricks.spark.csv')
    .options(header='true', inferSchema='true')
    .load(file1))

加载第二个文件:

schema = StructType(
    [StructField(x, LongType(), False) for x in ("c_id", "cn_id", "cn_value")])

df2 = (sqlContext.read 
    .format('com.databricks.spark.csv')
    .schema(schema)
    .options(header='false')
    .load(file2))

终于加入:

combined = df1.join(df2, df1["c_id"] == df2["c_id"])
combined.show()

## +------+----+------+-----+--------+
## |  c_id|name|  c_id|cn_id|cn_value|
## +------+----+------+-----+--------+
## |121210| abc|121210|    0|       1|
## |121211| pqr|121211|    0|       0|
## +------+----+------+-----+--------+

编辑:

使用 RDD,您可以执行以下操作:

file1_fields.join(file2_fields.map(lambda x: (x[0], x[1:])))

【讨论】:

  • 我们能否在不使用数据框的情况下实现它,例如使用 Pure Spark 转换和操作。通过制作键并将 c_id 替换为名称并聚合 cn_value...
  • 你可以,但是由于你的输入是表格的,它没有任何意义。
  • 谢谢,file1 中的数据以逗号分隔,包含多个字段,file2 中的数据也以逗号分隔,只有三个字段,但 file2 不是 csv。 file2 是 hdfs 部分文件,我必须映射和聚合这两个文件的结果。基于来自 file1 的 c_id 和来自 file2 的 c_id 然后用各自的名称替换 c_id 并基于来自 file2 的 c_id 和 cn_id 聚合结果。
  • file1_fields.join(file2_fields.map(lambda x: (x[0], x[1:]))) 这给了我像 (u'121210',(u'abc' , (u'0', 0))) 我们可以删除这些括号并得到像 121210,abc,0,0 这样的操作。请看
  • join 输出以下结构 (key, (value_left, value_right))。这是一个标准的 Python 元组,因此您可以根据需要对其进行整形。
猜你喜欢
  • 1970-01-01
  • 2017-09-29
  • 2017-08-10
  • 1970-01-01
  • 1970-01-01
  • 2016-12-23
  • 1970-01-01
  • 1970-01-01
  • 2017-08-19
相关资源
最近更新 更多