Class KafkaWriter
Write
Records to an Apache Kafka distributed messaging system.-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
Fields inherited from class com.northconcepts.datapipeline.core.DataEndpoint
lastRecord, PRODUCT, PRODUCT_VERSION, VENDOR, XML_INPUT_FACTORY_KEYFields inherited from class com.northconcepts.datapipeline.core.Endpoint
BUFFER_SIZE, captureElapsedTime, DEFAULT_READ_BUFFER_SIZEFields inherited from class com.northconcepts.datapipeline.core.DataObject
id, log, name, TIMESTAMP_FORMAT -
Constructor Summary
ConstructorsConstructorDescriptionKafkaWriter(Properties properties, String topic) Creates a writer that sends each record to the topic without a key; opening it puts the key/value serializers into the given properties. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.voidclose()Indicates that this endpoint has finished reading or writing.voidopen()Makes this endpoint ready for reading or writing.protected voidOverridden by subclasses to write the specified record to thisDataWriter.Methods inherited from class com.northconcepts.datapipeline.core.DataWriter
available, getNestedEndpoint, getNestedWriter, getRootEndpoint, getRootWriter, getWriter, writeMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, resetRecordCount, toStringMethods inherited from class com.northconcepts.datapipeline.core.Endpoint
addElapsedtime, assertClosed, assertNotOpened, assertOpened, finalize, getClosedOn, getDescription, getElapsedTime, getElapsedTimeAsString, getOpenedOn, getOpenElapsedTime, getOpenElapsedTimeAsString, getSelfTime, getSelfTimeAsString, getState, isCaptureElapsedTime, isClosed, isOpen, setCaptureElapsedTime, setDescription
-
Constructor Details
-
KafkaWriter
Creates a writer that sends each record to the topic without a key; opening it puts the key/value serializers into the given properties.
-
-
Method Details
-
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classIntegrationWriter- Throws:
DataException
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classDataEndpoint- Throws:
DataException
-
writeImpl
Description copied from class:DataWriterOverridden by subclasses to write the specified record to thisDataWriter.Contract for subclasses (see also
docs/authoring/DataWriter.md):- Do not call this method directly —
DataWriter.write(Record)is the template method that wraps exceptions, increments the record count, and attaches the offending record to thrownDataExceptions. - Do not call
DataEndpoint.incrementRecordCount()here;DataWriter.write(Record)already does. - Throw raw exceptions;
DataWriter.write(Record)wraps them viaexception(throwable).setRecord(record).
- Specified by:
writeImplin classDataWriter- Throws:
Throwable
- Do not call this method directly —
-
addExceptionProperties
Description copied from class:EndpointAdds this endpoint's current state to aDataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.- Overrides:
addExceptionPropertiesin classDataWriter
-