Class TypedCombineFnDelegate<InputT,AccumT,OutputT>
- java.lang.Object
-
- org.apache.beam.sdk.transforms.Combine.CombineFn<InputT,AccumT,OutputT>
-
- org.apache.beam.sdk.extensions.sql.TypedCombineFnDelegate<InputT,AccumT,OutputT>
-
- Type Parameters:
InputT- the type of inputAccumT- the type of accumulatorOutputT- the type of output
- All Implemented Interfaces:
java.io.Serializable,org.apache.beam.sdk.transforms.CombineFnBase.GlobalCombineFn<InputT,AccumT,OutputT>,org.apache.beam.sdk.transforms.display.HasDisplayData
public class TypedCombineFnDelegate<InputT,AccumT,OutputT> extends org.apache.beam.sdk.transforms.Combine.CombineFn<InputT,AccumT,OutputT>ACombine.CombineFndelegating all relevant calls to given delegate. This is used to create a type anonymous class for cases where the CombineFn is a generic class. The anonymous class can then be used in a UDAF as.registerUdaf("UDAF", new TypedCombineFnDelegate<>(genericCombineFn) {})- See Also:
- Serialized Form
-
-
Constructor Summary
Constructors Modifier Constructor Description protectedTypedCombineFnDelegate(org.apache.beam.sdk.transforms.Combine.CombineFn<InputT,AccumT,OutputT> delegate)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description AccumTaddInput(AccumT mutableAccumulator, InputT input)OutputTapply(java.lang.Iterable<? extends InputT> inputs)AccumTcompact(AccumT accumulator)AccumTcreateAccumulator()OutputTdefaultValue()OutputTextractOutput(AccumT accumulator)java.lang.reflect.TypeVariable<?>getAccumTVariable()org.apache.beam.sdk.coders.Coder<AccumT>getAccumulatorCoder(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<InputT> inputCoder)org.apache.beam.sdk.coders.Coder<OutputT>getDefaultOutputCoder(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<InputT> inputCoder)java.lang.StringgetIncompatibleGlobalWindowErrorMessage()java.lang.reflect.TypeVariable<?>getInputTVariable()org.apache.beam.sdk.values.TypeDescriptor<InputT>getInputType()java.lang.reflect.TypeVariable<?>getOutputTVariable()org.apache.beam.sdk.values.TypeDescriptor<OutputT>getOutputType()AccumTmergeAccumulators(java.lang.Iterable<AccumT> accumulators)voidpopulateDisplayData(org.apache.beam.sdk.transforms.display.DisplayData.Builder builder)
-
-
-
Method Detail
-
getOutputType
public org.apache.beam.sdk.values.TypeDescriptor<OutputT> getOutputType()
-
getInputType
public org.apache.beam.sdk.values.TypeDescriptor<InputT> getInputType()
-
createAccumulator
public AccumT createAccumulator()
-
defaultValue
public OutputT defaultValue()
-
getAccumulatorCoder
public org.apache.beam.sdk.coders.Coder<AccumT> getAccumulatorCoder(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<InputT> inputCoder) throws org.apache.beam.sdk.coders.CannotProvideCoderException
-
getDefaultOutputCoder
public org.apache.beam.sdk.coders.Coder<OutputT> getDefaultOutputCoder(org.apache.beam.sdk.coders.CoderRegistry registry, org.apache.beam.sdk.coders.Coder<InputT> inputCoder) throws org.apache.beam.sdk.coders.CannotProvideCoderException
-
getIncompatibleGlobalWindowErrorMessage
public java.lang.String getIncompatibleGlobalWindowErrorMessage()
-
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 builder)
- Specified by:
populateDisplayDatain interfaceorg.apache.beam.sdk.transforms.display.HasDisplayData
-
-