Class CovarianceFn<T extends java.lang.Number>

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

    @Internal
    public class CovarianceFn<T extends java.lang.Number>
    extends org.apache.beam.sdk.transforms.Combine.CombineFn<org.apache.beam.sdk.values.Row,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator,​T>
    Combine.CombineFn for Covariance on Number types.

    Calculates Population Covariance and Sample Covariance using incremental formulas described in http://en.wikipedia.org/wiki/Algorithms_for_calculating_variance, presumably by Pébay, Philippe (2008), in "Formulas for Robust, One-Pass Parallel Computation of Covariances and Arbitrary-Order Statistical Moments".

    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.CovarianceAccumulator addInput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator currentVariance, org.apache.beam.sdk.values.Row rawInput)  
      org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator createAccumulator()  
      T extractOutput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator accumulator)  
      java.lang.reflect.TypeVariable<?> getAccumTVariable()  
      org.apache.beam.sdk.coders.Coder<org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator> getAccumulatorCoder​(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<org.apache.beam.sdk.values.Row> 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.CovarianceAccumulator mergeAccumulators​(java.lang.Iterable<org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator> covariances)  
      static CovarianceFn newPopulation​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)  
      static <V extends java.lang.Number>
      CovarianceFn
      newPopulation​(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,​V> decimalConverter)  
      static CovarianceFn newSample​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)  
      static <V extends java.lang.Number>
      CovarianceFn
      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 CovarianceFn newPopulation​(org.apache.beam.sdk.schemas.Schema.TypeName typeName)
      • newPopulation

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

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

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

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

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

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

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

        public T extractOutput​(org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator accumulator)
        Specified by:
        extractOutput in class org.apache.beam.sdk.transforms.Combine.CombineFn<org.apache.beam.sdk.values.Row,​org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator,​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