Class KafkaTopicSink<T>
java.lang.Object
org.apache.wayang.core.plan.wayangplan.OperatorBase
org.apache.wayang.core.plan.wayangplan.UnarySink<T>
org.apache.wayang.basic.operators.KafkaTopicSink<T>
- All Implemented Interfaces:
Serializable
,ActualOperator
,ElementaryOperator
,Operator
- Direct Known Subclasses:
JavaKafkaTopicSink
,SparkKafkaTopicSink
This
UnarySink
writes all incoming data quanta to a single Kafka topic.- See Also:
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
OperatorBase.GsonSerializer
-
Field Summary
FieldsFields inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
inputSlots, outputSlots, STANDARD_OPERATOR_ARGS
Fields inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
FIRST_EPOCH
-
Constructor Summary
ConstructorsConstructorDescriptionKafkaTopicSink
(String topicName, Class<T> typeClass) Creates a new instance with default formatting.KafkaTopicSink
(String topicName, FunctionDescriptor.SerializableFunction<T, String> formattingFunction, Class<T> typeClass) Creates a new instance.KafkaTopicSink
(String topicName, TransformationDescriptor<T, String> formattingDescriptor) Creates a new instance.KafkaTopicSink
(KafkaTopicSink<T> that) Creates a copied instance. -
Method Summary
Modifier and TypeMethodDescriptionstatic Properties
getProducer
(Properties props) void
initProducer
(KafkaTopicSink<T> kts) static Properties
loadConfig
(String propertiesFilePath) Load properties from a properties file or alternatively use the default properties with some sensitive values from environment variables.Methods inherited from class org.apache.wayang.core.plan.wayangplan.OperatorBase
accept, addBroadcastInput, addTargetPlatform, at, collectMappedInputSlots, collectMappedOutputSlots, copy, createCopy, getAllInputs, getAllOutputs, getCardinalityEstimator, getContainer, getEpoch, getName, getOriginal, getSimpleClassName, getTargetPlatforms, isAuxiliary, isSupportingBroadcastInputs, propagateInputCardinality, propagateOutputCardinality, setAuxiliary, setCardinalityEstimator, setContainer, setEpoch, setName, toString
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, wait
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.ActualOperator
accept
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.ElementaryOperator
createCardinalityEstimator, getCardinalityEstimator, isAuxiliary, setAuxiliary, setCardinalityEstimator
Methods inherited from interface org.apache.wayang.core.plan.wayangplan.Operator
addBroadcastInput, addTargetPlatform, broadcastTo, broadcastTo, collectMappedInputSlots, collectMappedOutputSlots, connectTo, connectTo, getAllInputs, getAllOutputs, getCardinalityPusher, getContainer, getEffectiveOccupant, getEffectiveOccupant, getEpoch, getEstimationContextProperties, getForwards, getInnermostLoop, getInput, getInput, getLoopStack, getName, getNumBroadcastInputs, getNumInputs, getNumOutputs, getNumRegularInputs, getOuterInputSlot, getOutermostInputSlot, getOutermostOutputSlots, getOutput, getOutput, getParent, getTargetPlatforms, isAlternative, isConversion, isElementary, isExecutionOperator, isFeedbackInput, isFeedforwardOutput, isLoopHead, isLoopSubplan, isOwnerOf, isReading, isSink, isSource, isSubplan, isSupportingBroadcastInputs, isUnconnected, propagateInputCardinality, propagateOutputCardinality, propagateOutputCardinality, replaceWith, setContainer, setEpoch, setInput, setName, setOutput
-
Field Details
-
topicName
-
formattingDescriptor
-
-
Constructor Details
-
KafkaTopicSink
public KafkaTopicSink() -
KafkaTopicSink
Creates a new instance with default formatting.- Parameters:
topicName
- Name of Kafka topic that should be written totypeClass
-Class
of incoming data quanta
-
KafkaTopicSink
public KafkaTopicSink(String topicName, FunctionDescriptor.SerializableFunction<T, String> formattingFunction, Class<T> typeClass) Creates a new instance. -
KafkaTopicSink
Creates a new instance.- Parameters:
topicName
- Name of Kafka topic that should be written toformattingDescriptor
- formats incoming data quanta to aString
representation
-
KafkaTopicSink
Creates a copied instance.- Parameters:
that
- should be copied
-
-
Method Details
-
initProducer
-
getProducer
-
getProducer
-
loadConfig
Load properties from a properties file or alternatively use the default properties with some sensitive values from environment variables.- Parameters:
propertiesFilePath
- - File path or null.- Returns:
- Properties object
-
getDefaultProperties
-