【问题标题】:How to compute standard deviation in Apache Beam如何计算 Apache Beam 中的标准差
【发布时间】:2018-08-13 21:38:56
【问题描述】:

我是 Apache Beam 新手,我想计算大型数据集的均值和标准偏差。

给定一个“A,B”形式的 .csv 文件,其中 A,B 是整数,这基本上就是我所拥有的。

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.io.textio import ReadFromText

class Split(beam.DoFn):
    def process(self, element):
        A, B = element.split(',')
        return [('A', A), ('B', B)]

with beam.Pipeline(options=PipelineOptions()) as p:
     # parse the rows
     rows = (p
             | ReadFromText('data.csv')
             | beam.ParDo(Split()))

     # calculate the mean
     avgs = (rows
             | beam.CombinePerKey(
                 beam.combiners.MeanCombineFn()))

     # calculate the stdv per key
     # ???

     std >> beam.io.WriteToText('std.out')

我想做这样的事情:

class SquaredDiff(beam.DoFn):
    def process(self, element):
        A = element[0][1]
        B = element[1][1]
        return [('A', A - avgs[0]), ('B', B - avgs[1])]

stdv = (rows
        | beam.ParDo(SquaredDiff())
        | beam.CombinePerKey(
            beam.combiners.MeanCombineFn()))

什么的,但我不知道怎么做。

【问题讨论】:

    标签: python apache-beam


    【解决方案1】:

    编写您自己的组合器。这将起作用:

    class MeanStddev(beam.CombineFn):
      def create_accumulator(self):
        return (0.0, 0.0, 0) # x, x^2, count
    
      def add_input(self, sum_count, input):
        (sum, sumsq, count) = sum_count
        return sum + input, sumsq + input*input, count + 1
    
      def merge_accumulators(self, accumulators):
        sums, sumsqs, counts = zip(*accumulators)
        return sum(sums), sum(sumsqs), sum(counts)
    
      def extract_output(self, sum_count):
        (sum, sumsq, count) = sum_count
        if count:
          mean = sum / count
          variance = (sumsq / count) - mean*mean
          # -ve value could happen due to rounding
          stddev = np.sqrt(variance) if variance > 0 else 0
          return {
            'mean': mean,
            'variance': variance,
            'stddev': stddev,
            'count': count
          }
        else:
          return {
            'mean': float('NaN'),
            'variance': float('NaN'),
            'stddev': float('NaN'),
            'count': 0
          }
    

    这会将方差计算为 E(x^2) - E(x)*E(x),因此您只需通过数据一次。这就是你将如何使用上面的组合器:

    [1.3, 3.0, 4.2] | beam.CombineGlobally(MeanStddev())
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-01-21
      • 2019-09-22
      • 2016-04-14
      • 2020-01-28
      • 2020-02-28
      • 2013-06-03
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多