Class CovarianceFn<T extends java.lang.Number>
- java.lang.Object
-
- org.apache.beam.sdk.transforms.Combine.CombineFn<org.apache.beam.sdk.values.Row,org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator,T>
-
- org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceFn<T>
-
- 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.CombineFnfor Covariance onNumbertypes.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.CovarianceAccumulatoraddInput(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.CovarianceAccumulatorcreateAccumulator()TextractOutput(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.StringgetIncompatibleGlobalWindowErrorMessage()java.lang.reflect.TypeVariable<?>getInputTVariable()java.lang.reflect.TypeVariable<?>getOutputTVariable()org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulatormergeAccumulators(java.lang.Iterable<org.apache.beam.sdk.extensions.sql.impl.transform.agg.CovarianceAccumulator> covariances)static CovarianceFnnewPopulation(org.apache.beam.sdk.schemas.Schema.TypeName typeName)static <V extends java.lang.Number>
CovarianceFnnewPopulation(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,V> decimalConverter)static CovarianceFnnewSample(org.apache.beam.sdk.schemas.Schema.TypeName typeName)static <V extends java.lang.Number>
CovarianceFnnewSample(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,V> decimalConverter)voidpopulateDisplayData(org.apache.beam.sdk.transforms.display.DisplayData.Builder arg0)
-
-
-
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:
createAccumulatorin classorg.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:
addInputin classorg.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:
mergeAccumulatorsin classorg.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:
getAccumulatorCoderin interfaceorg.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:
extractOutputin classorg.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:
getDefaultOutputCoderin interfaceorg.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:
getIncompatibleGlobalWindowErrorMessagein interfaceorg.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:
populateDisplayDatain interfaceorg.apache.beam.sdk.transforms.display.HasDisplayData
-
-