【问题标题】:Sampling a large distributed data set using pyspark / spark使用 pyspark / spark 对大型分布式数据集进行采样
【发布时间】: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


【解决方案1】:

尝试改用textFile.sample(false,fraction,seed)。 takeSample 通常会很慢,因为它 calls count() on the RDD。它需要这样做,因为否则它不会从每个分区中平均获取,基本上它使用计数以及您要求的样本大小来计算分数并在内部调用sample。 sample 很快,因为它只使用一个随机布尔生成器,该生成器在百分比的时间内返回 true fraction,因此不需要调用 count。

此外,我认为这不会发生在您身上,但如果返回的样本量不够大,它会再次调用sample,这显然会减慢速度。由于您应该对数据的大小有所了解,我建议您调用 sample 然后将样本缩减到自己的大小,因为您比 spark 更了解您的数据。

【讨论】:

  • 这有点奇怪。计数并不是一个缓慢的操作 - 它比 takeSample 快了约 2 个数量级,这表明这不是核心问题。
【解决方案2】:

使用 sample 而不是 takeSample 似乎使事情变得相当快:

textFile.sample(False, .0001, 12345)

这样做的问题是,除非您对数据集中的行数有一个粗略的了解,否则很难知道要选择的正确分数。

【讨论】:

    【解决方案3】:

    PySpark 中不同类型的样本

    随机抽样 % 的数据,无论是否替换

    import pyspark.sql.functions as F
    #Randomly sample 50% of the data without replacement
    sample1 = df.sample(False, 0.5, seed=0)
    
    #Randomly sample 50% of the data with replacement
    sample1 = df.sample(True, 0.5, seed=0)
    
    #Take another sample exlcuding records from previous sample using Anti Join
    sample2 = df.join(sample1, on='ID', how='left_anti').sample(False, 0.5, seed=0)
    
    #Take another sample exlcuding records from previous sample using Where
    sample1_ids = [row['ID'] for row in sample1.ID]
    sample2 = df.where(~F.col('ID').isin(sample1_ids)).sample(False, 0.5, seed=0)
    
    #Generate a startfied sample of the data across column(s)
    #Sampling is probabilistic and thus cannot guarantee an exact number of rows
    fractions = {
            'NJ': 0.5, #Take about 50% of records where state = NJ
        'NY': 0.25, #Take about 25% of records where state = NY
        'VA': 0.1, #Take about 10% of records where state = VA
    }
    stratified_sample = df.sampleBy(F.col('state'), fractions, seed=0)
    

    【讨论】:

      猜你喜欢
      • 2018-05-18
      • 1970-01-01
      • 1970-01-01
      • 2016-01-10
      • 2018-10-29
      • 2017-05-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多