Class VarianceFn<T extends java.lang.Number>
- java.lang.Object
-
- org.apache.beam.sdk.transforms.Combine.CombineFn<T,org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator,T>
-
- org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceFn<T>
-
- 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.CombineFnfor Variance onNumbertypes.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.VarianceAccumulatoraddInput(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator currentVariance, T rawInput)org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulatorcreateAccumulator()TextractOutput(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.StringgetIncompatibleGlobalWindowErrorMessage()java.lang.reflect.TypeVariable<?>getInputTVariable()java.lang.reflect.TypeVariable<?>getOutputTVariable()org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulatormergeAccumulators(java.lang.Iterable<org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator> variances)static VarianceFnnewPopulation(org.apache.beam.sdk.schemas.Schema.TypeName typeName)static <V extends java.lang.Number>
VarianceFnnewPopulation(org.apache.beam.sdk.transforms.SerializableFunction<java.math.BigDecimal,V> decimalConverter)static VarianceFnnewSample(org.apache.beam.sdk.schemas.Schema.TypeName typeName)static <V extends java.lang.Number>
VarianceFnnewSample(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 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()
-
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)
-
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)
-
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)
-
extractOutput
public T extractOutput(org.apache.beam.sdk.extensions.sql.impl.transform.agg.VarianceAccumulator accumulator)
-
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
-
-