OUT
- the output type of the source.public class KafkaSource<OUT> extends Object implements Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>, ResultTypeQueryable<OUT>
KafkaSourceBuilder
to construct a KafkaSource
. The following example shows how to create a KafkaSource emitting records of
String
type.
KafkaSource<String> source = KafkaSource
.<String>builder()
.setBootstrapServers(KafkaSourceTestEnv.brokerConnectionStrings)
.setGroupId("MyGroup")
.setTopics(Arrays.asList(TOPIC1, TOPIC2))
.setDeserializer(new TestingKafkaRecordDeserializer())
.setStartingOffsets(OffsetsInitializer.earliest())
.build();
See KafkaSourceBuilder
for more details.
public static <OUT> KafkaSourceBuilder<OUT> builder()
KafkaSource
.public Boundedness getBoundedness()
Source
getBoundedness
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
public SourceReader<OUT,KafkaPartitionSplit> createReader(SourceReaderContext readerContext)
Source
createReader
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
readerContext
- The context
for the source reader.public SplitEnumerator<KafkaPartitionSplit,KafkaSourceEnumState> createEnumerator(SplitEnumeratorContext<KafkaPartitionSplit> enumContext)
Source
createEnumerator
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
enumContext
- The context
for the split enumerator.public SplitEnumerator<KafkaPartitionSplit,KafkaSourceEnumState> restoreEnumerator(SplitEnumeratorContext<KafkaPartitionSplit> enumContext, KafkaSourceEnumState checkpoint) throws IOException
Source
restoreEnumerator
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
enumContext
- The context
for the restored split
enumerator.checkpoint
- The checkpoint to restore the SplitEnumerator from.IOException
public SimpleVersionedSerializer<KafkaPartitionSplit> getSplitSerializer()
Source
getSplitSerializer
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
public SimpleVersionedSerializer<KafkaSourceEnumState> getEnumeratorCheckpointSerializer()
Source
SplitEnumerator
checkpoint. The serializer is used for
the result of the SplitEnumerator.snapshotState()
method.getEnumeratorCheckpointSerializer
in interface Source<OUT,KafkaPartitionSplit,KafkaSourceEnumState>
public TypeInformation<OUT> getProducedType()
ResultTypeQueryable
TypeInformation
) produced by this function or input format.getProducedType
in interface ResultTypeQueryable<OUT>
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.