Modifier and Type | Class and Description |
---|---|
class |
GuavaFlinkConnectorRateLimiter
An implementation of
FlinkConnectorRateLimiter that uses Guava's RateLimiter for rate
limiting. |
Modifier and Type | Field and Description |
---|---|
protected FlinkConnectorRateLimiter |
PubSubSource.rateLimiter |
Modifier and Type | Method and Description |
---|---|
FlinkConnectorRateLimiter |
FlinkKafkaConsumer010.getRateLimiter() |
Modifier and Type | Method and Description |
---|---|
void |
FlinkKafkaConsumer010.setRateLimiter(FlinkConnectorRateLimiter kafkaRateLimiter)
Set a rate limiter to ratelimit bytes read from Kafka.
|
Constructor and Description |
---|
Kafka010Fetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics,
FlinkConnectorRateLimiter rateLimiter) |
KafkaConsumerThread(org.slf4j.Logger log,
Handover handover,
Properties kafkaProperties,
ClosableBlockingQueue<KafkaTopicPartitionState<T,org.apache.kafka.common.TopicPartition>> unassignedPartitionsQueue,
String threadName,
long pollTimeout,
boolean useMetrics,
MetricGroup consumerMetricGroup,
MetricGroup subtaskMetricGroup,
FlinkConnectorRateLimiter rateLimiter) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.