【发布时间】:2014-09-08 11:23:56
【问题描述】:
我在 hdfs 中有一个文件,它分布在集群中的节点上。
我正在尝试从此文件中随机抽取 10 行样本。
在 pyspark shell 中,我使用以下方法将文件读入 RDD:
>>> textFile = sc.textFile("/user/data/myfiles/*")
然后我想简单地取样...关于 Spark 很酷的一点是有像 takeSample 这样的命令,不幸的是我认为我做错了什么,因为以下需要很长时间:
>>> textFile.takeSample(False, 10, 12345)
所以我尝试在每个节点上创建一个分区,然后使用以下命令指示每个节点对该分区进行采样:
>>> textFile.partitionBy(4).mapPartitions(lambda blockOfLines: blockOfLines.takeSample(False, 10, 1234)).first()
但这会产生错误ValueError: too many values to unpack:
org.apache.spark.api.python.PythonException: Traceback (most recent call last):
File "/opt/cloudera/parcels/CDH-5.0.2-1.cdh5.0.2.p0.13/lib/spark/python/pyspark/worker.py", line 77, in main
serializer.dump_stream(func(split_index, iterator), outfile)
File "/opt/cloudera/parcels/CDH-5.0.2-1.cdh5.0.2.p0.13/lib/spark/python/pyspark/serializers.py", line 117, in dump_stream
for obj in iterator:
File "/opt/cloudera/parcels/CDH-5.0.2-1.cdh5.0.2.p0.13/lib/spark/python/pyspark/rdd.py", line 821, in add_shuffle_key
for (k, v) in iterator:
ValueError: too many values to unpack
如何使用 spark 或 pyspark 从大型分布式数据集中采样 10 行?
【问题讨论】:
-
我认为这不是 spark 的问题,请参阅 stackoverflow.com/questions/7053551/…
-
@aaronman 你是正确的,因为“太多的值”错误绝对是一个 python 错误。我将添加有关错误消息的更多详细信息。我的预感是我的 pyspark 代码有问题 - 您能否在 spark 设置上成功运行此代码?
-
我只真正使用了 scala spark API,我认为 scala 的函数式风格非常适合 Mapreduce
-
@aaronman 我愿意接受 scala 解决方案!
-
@samthebest - 我不知道我是否在这里遗漏了一些东西,python 和 scala 都是函数式语言,而 spark 有 Python 和 Scala API。这不仅仅是一个偏好问题吗?
标签: hadoop apache-spark