Package | Description |
---|---|
org.apache.flink.streaming.connectors.kafka | |
org.apache.flink.streaming.connectors.kafka.config | |
org.apache.flink.streaming.connectors.kafka.shuffle |
Modifier and Type | Method and Description |
---|---|
protected static void |
FlinkKafkaConsumerBase.adjustAutoCommitConfig(Properties properties,
OffsetCommitMode offsetCommitMode)
Make sure that auto commit is disabled when our offset commit mode is ON_CHECKPOINTS.
|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> subscribedPartitionsToStartOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup kafkaMetricGroup,
boolean useMetrics)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the data, and
emits it into the data streams.
|
Modifier and Type | Method and Description |
---|---|
static OffsetCommitMode |
OffsetCommitModes.fromConfiguration(boolean enableAutoCommit,
boolean enableCommitOnCheckpoint,
boolean enableCheckpointing)
Determine the offset commit mode using several configuration values.
|
static OffsetCommitMode |
OffsetCommitMode.valueOf(String name)
Returns the enum constant of this type with the specified name.
|
static OffsetCommitMode[] |
OffsetCommitMode.values()
Returns an array containing the constants of this enum type, in
the order they are declared.
|
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaShuffleConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.