Class VarianceFn<T extends java.lang.Number>

  • All Implemented Interfaces:
    java.io.Serializable, org.apache.beam.sdk.transforms.CombineFnBase.GlobalCombineFn<T,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T>, org.apache.beam.sdk.transforms.display.HasDisplayData

    @Internal
    public class VarianceFn<T extends java.lang.Number>
    extends org.apache.beam.sdk.transforms.Combine.CombineFn<T,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T>
    Combine.CombineFn for Variance on Number types.

    Calculates Population Variance and Sample Variance using incremental formulas described, for example, by Chan, Golub, and LeVeque in "Algorithms for computing the sample variance: analysis and recommendations", The American Statistician, 37 (1983) pp. 242--247.

    If variance is defined like this:

    • Input elements: (x[1], ... , x[n])
    • Sum of elements: {sum(x) = x[1] + ... + x[n]}
    • Average of all elements in the input: mean(x) = sum(x) / n
    • Deviation of ith element from the current mean: deviation(x, i) = x[i] - mean(n)
    • Variance: variance(x) = deviation(x, 1)^2 + ... + deviation(x, n)^2

    Then variance of combined input of 2 samples (x[1], ... , x[n]) and (y[1], ... , y[m]) is calculated using this formula:

    • variance(concat(x,y)) = variance(x) + variance(y) + increment, where:
    • increment = m/(n(m+n)) * (n/m * sum(x) - sum(y))^2

    This is also applicable for a single element increment, assuming that variance of a single element input is zero

    To implement the above formula we keep track of the current variation, sum, and count of elements, and then use the formula whenever new element comes or we need to merge variances for 2 samples.

    See Also:
    Serialized Form
    • Method Summary

      All Methods Static Methods Instance Methods Concrete Methods 
      Modifier and Type Method Description
      org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator addInput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator currentVariance, T rawInput)  
      org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator createAccumulator()  
      T extractOutput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator accumulator)  
      java.lang.reflect.TypeVariable<?> getAccumTVariable()  
      org.apache.beam.sdk.coders.Coder<org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator> getAccumulatorCoder​(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<T> inputCoder)  
      org.apache.beam.sdk.coders.Coder<OutputT> getDefaultOutputCoder​(org.apache.beam.sdk.coders.CoderRegistry arg0, org.apache.beam.sdk.coders.Coder<@UnknownKeyFor @NonNull @Initialized InputT> arg1)  
      java.lang.String getIncompatibleGlobalWindowErrorMessage()  
      java.lang.reflect.TypeVariable<?> getInputTVariable()  
      java.lang.reflect.TypeVariable<?> getOutputTVariable()  
      org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator mergeAccumulators​(java.lang.Iterable<org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator> variances)  
      static VarianceFn newPopulation​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)  
      static <V extends java.lang.Number>
      VarianceFn
      newPopulation​(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,​V> decimalConverter)  
      static VarianceFn newSample​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)  
      static <V extends java.lang.Number>
      VarianceFn
      newSample​(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,​V> decimalConverter)  
      void populateDisplayData​(org.apache.beam.sdk.transforms.display.DisplayData.Builder arg0)  
      • Methods inherited from class org.apache.beam.sdk.transforms.Combine.CombineFn

        apply, compact, defaultValue, getInputType, getOutputType
      • Methods inherited from class java.lang.Object

        clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
    • Method Detail

      • newPopulation

        public static VarianceFn newPopulation​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)
      • newPopulation

        public static <V extends java.lang.Number> VarianceFn newPopulation​(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,​V> decimalConverter)
      • newSample

        public static VarianceFn newSample​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)
      • newSample

        public static <V extends java.lang.Number> VarianceFn newSample​(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,​V> decimalConverter)
      • createAccumulator

        public org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator createAccumulator()
        Specified by:
        createAccumulator in class org.apache.beam.sdk.transforms.Combine.CombineFn<T extends java.lang.Number,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T extends java.lang.Number>
      • addInput

        public org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator addInput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator currentVariance,
                                                                                                  T rawInput)
        Specified by:
        addInput in class org.apache.beam.sdk.transforms.Combine.CombineFn<T extends java.lang.Number,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T extends java.lang.Number>
      • mergeAccumulators

        public org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator mergeAccumulators​(java.lang.Iterable<org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator> variances)
        Specified by:
        mergeAccumulators in class org.apache.beam.sdk.transforms.Combine.CombineFn<T extends java.lang.Number,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T extends java.lang.Number>
      • getAccumulatorCoder

        public org.apache.beam.sdk.coders.Coder<org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator> getAccumulatorCoder​(org.apache.beam.sdk.coders.CoderRegistry registry,
                                                                                                                                               org.apache.beam.sdk.coders.Coder<T> inputCoder)
        Specified by:
        getAccumulatorCoder in interface org.apache.beam.sdk.transforms.CombineFnBase.GlobalCombineFn<T extends java.lang.Number,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T extends java.lang.Number>
      • extractOutput

        public T extractOutput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator accumulator)
        Specified by:
        extractOutput in class org.apache.beam.sdk.transforms.Combine.CombineFn<T extends java.lang.Number,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,​T extends java.lang.Number>
      • getDefaultOutputCoder

        public org.apache.beam.sdk.coders.Coder<OutputT> getDefaultOutputCoder​(org.apache.beam.sdk.coders.CoderRegistry arg0,
                                                                               org.apache.beam.sdk.coders.Coder<@UnknownKeyFor @NonNull @Initialized InputT> arg1)
                                                                        throws org.apache.beam.sdk.coders.CannotProvideCoderException
        Specified by:
        getDefaultOutputCoder in interface org.apache.beam.sdk.transforms.CombineFnBase.GlobalCombineFn<InputT extends java.lang.Object,​AccumT extends java.lang.Object,​OutputT extends java.lang.Object>
        Throws:
        org.apache.beam.sdk.coders.CannotProvideCoderException
      • getIncompatibleGlobalWindowErrorMessage

        public java.lang.String getIncompatibleGlobalWindowErrorMessage()
        Specified by:
        getIncompatibleGlobalWindowErrorMessage in interface org.apache.beam.sdk.transforms.CombineFnBase.GlobalCombineFn<InputT extends java.lang.Object,​AccumT extends java.lang.Object,​OutputT extends java.lang.Object>
      • getInputTVariable

        public java.lang.reflect.TypeVariable<?> getInputTVariable()
      • getAccumTVariable

        public java.lang.reflect.TypeVariable<?> getAccumTVariable()
      • getOutputTVariable

        public java.lang.reflect.TypeVariable<?> getOutputTVariable()
      • populateDisplayData

        public void populateDisplayData​(org.apache.beam.sdk.transforms.display.DisplayData.Builder arg0)
        Specified by:
        populateDisplayData in interface org.apache.beam.sdk.transforms.display.HasDisplayData