Class BeamProducer
java.lang.Object
org.apache.beam.sdk.transforms.PTransform<@NonNull org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,@NonNull org.apache.beam.sdk.values.PDone>
com.logicalclocks.hsfs.beam.engine.BeamProducer
- All Implemented Interfaces:
Serializable,org.apache.beam.sdk.transforms.display.HasDisplayData
public class BeamProducer
extends org.apache.beam.sdk.transforms.PTransform<@NonNull org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,@NonNull org.apache.beam.sdk.values.PDone>
- See Also:
-
Field Summary
Fields inherited from class org.apache.beam.sdk.transforms.PTransform
name, resourceHints -
Constructor Summary
ConstructorsConstructorDescriptionBeamProducer(String topic, Map<String, String> properties, org.apache.avro.Schema schema, org.apache.avro.Schema encodedSchema, Map<String, org.apache.avro.Schema> deserializedComplexFeatureSchemas, List<String> primaryKeys, StreamFeatureGroup streamFeatureGroup) -
Method Summary
Modifier and TypeMethodDescriptionorg.apache.beam.sdk.values.PDoneexpand(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input) Methods inherited from class org.apache.beam.sdk.transforms.PTransform
compose, compose, getAdditionalInputs, getDefaultOutputCoder, getDefaultOutputCoder, getDefaultOutputCoder, getKindString, getName, getResourceHints, populateDisplayData, setResourceHints, toString, validate, validate
-
Constructor Details
-
BeamProducer
public BeamProducer(String topic, Map<String, String> properties, org.apache.avro.Schema schema, org.apache.avro.Schema encodedSchema, Map<String, throws FeatureStoreException, IOExceptionorg.apache.avro.Schema> deserializedComplexFeatureSchemas, List<String> primaryKeys, StreamFeatureGroup streamFeatureGroup) - Throws:
FeatureStoreExceptionIOException
-
-
Method Details
-
expand
public org.apache.beam.sdk.values.PDone expand(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input) - Specified by:
expandin classorg.apache.beam.sdk.transforms.PTransform<@NonNull org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,@NonNull org.apache.beam.sdk.values.PDone>
-