|
Change Parameter Type subscribedPartitionStates : KafkaTopicPartitionState<TopicPartition>[] to unassignedPartitionsQueue : ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>> in method public KafkaConsumerThread(log Logger, handover Handover, kafkaProperties Properties, unassignedPartitionsQueue ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>>, kafkaMetricGroup MetricGroup, consumerCallBridge KafkaConsumerCallBridge, threadName String, pollTimeout long, useMetrics boolean) in class org.apache.flink.streaming.connectors.kafka.internal.KafkaConsumerThread |
From |
To |
|
Change Parameter Type partitionStates : KafkaTopicPartitionState<?>[] to partitionStates : List<KafkaTopicPartitionState<TopicAndPartition>> in method package PeriodicOffsetCommitter(offsetHandler ZookeeperOffsetHandler, partitionStates List<KafkaTopicPartitionState<TopicAndPartition>>, errorHandler ExceptionProxy, commitInterval long) in class org.apache.flink.streaming.connectors.kafka.internals.PeriodicOffsetCommitter |
From |
To |
|
Change Parameter Type partitions : KafkaTopicPartitionState<TopicPartition>[] to partitions : List<KafkaTopicPartitionState<TopicPartition>> in method private convertKafkaPartitions(partitions List<KafkaTopicPartitionState<TopicPartition>>) : List<TopicPartition> in class org.apache.flink.streaming.connectors.kafka.internal.KafkaConsumerThread |
From |
To |
|
Change Parameter Type allPartitions : KafkaTopicPartitionStateWithPeriodicWatermarks<?,?>[] to allPartitions : List<KafkaTopicPartitionState<KPH>> in method package PeriodicWatermarkEmitter(allPartitions List<KafkaTopicPartitionState<KPH>>, emitter SourceContext<?>, timerService ProcessingTimeService, autoWatermarkInterval long) in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher.PeriodicWatermarkEmitter |
From |
To |
|
Change Variable Type offsets : HashMap<KafkaTopicPartition,Long> to offsets : Map<KafkaTopicPartition,Long> in method public notifyCheckpointComplete(checkpointId long) : void in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Variable Type kafkaTopicPartition : KafkaTopicPartition to restoredStateEntry : Map.Entry<KafkaTopicPartition,Long> in method public open(configuration Configuration) : void in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Variable Type newPartitions : List<KafkaTopicPartition> to clone : List<List<KafkaTopicPartition>> in method private deepClone(toClone List<List<KafkaTopicPartition>>) : List<List<KafkaTopicPartition>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Change Variable Type partition : KafkaTopicPartitionState<?> to partition : KafkaTopicPartitionState<KPH> in method public snapshotCurrentState() : HashMap<KafkaTopicPartition,Long> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Change Variable Type partitions : KafkaTopicPartitionState<TopicPartition>[] to partitions : List<KafkaTopicPartitionState<TopicPartition>> in method public commitInternalOffsetsToKafka(offsets Map<KafkaTopicPartition,Long>) : void in class org.apache.flink.streaming.connectors.kafka.internal.Kafka09Fetcher |
From |
To |
|
Change Variable Type ktp : KafkaTopicPartitionState<?> to ktp : KafkaTopicPartitionState<KPH> in method protected addOffsetStateGauge(metricGroup MetricGroup) : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Change Variable Type state : KafkaTopicPartitionStateWithPeriodicWatermarks<?,?> to state : KafkaTopicPartitionState<?> in method public onProcessingTime(timestamp long) : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher.PeriodicWatermarkEmitter |
From |
To |
|
Rename Variable e : Exception to expected : Exception in method public testAllBoostrapServerHostsAreInvalid() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumer08Test |
From |
To |
|
Rename Variable inPartitions : List<KafkaTopicPartition> to mockGetAllPartitionsForTopicsReturn : List<KafkaTopicPartition> in method public testPartitionsEqualConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable inPartitions : List<KafkaTopicPartition> to mockGetAllPartitionsForTopicsReturn : List<KafkaTopicPartition> in method public testPartitionsFewerThanConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable partition : Map.Entry<KafkaTopicPartition,Long> to partitionEntry : Map.Entry<KafkaTopicPartition,Long> in method private createPartitionStateHolders(partitionsToInitialOffsets Map<KafkaTopicPartition,Long>, timestampWatermarkMode int, watermarksPeriodic SerializedValue<AssignerWithPeriodicWatermarks<T>>, watermarksPunctuated SerializedValue<AssignerWithPunctuatedWatermarks<T>>, userCodeClassLoader ClassLoader) : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Rename Variable subscribedPartitions : List<KafkaTopicPartition> to initialDiscovery : List<KafkaTopicPartition> in method public testMultiplePartitionsPerConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable subscribedPartitions : List<KafkaTopicPartition> to initialDiscovery : List<KafkaTopicPartition> in method public testPartitionsEqualConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable subscribedPartitions : List<KafkaTopicPartition> to initialDiscovery : List<KafkaTopicPartition> in method public testPartitionsFewerThanConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable kafkaTopicPartition : KafkaTopicPartition to restoredStateEntry : Map.Entry<KafkaTopicPartition,Long> in method public open(configuration Configuration) : void in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Rename Variable newPartitions : List<KafkaTopicPartition> to clone : List<List<KafkaTopicPartition>> in method private deepClone(toClone List<List<KafkaTopicPartition>>) : List<List<KafkaTopicPartition>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Variable partitions : List<KafkaTopicPartition> to mockGetAllPartitionsForTopicsReturn : List<KafkaTopicPartition> in method public testMultiplePartitionsPerConsumersFixedPartitions() : void in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Change Return Type KafkaTopicPartitionState<KPH>[] to List<KafkaTopicPartitionState<KPH>> in method protected subscribedPartitionStates() : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Change Return Type HashMap<KafkaTopicPartition,Long> to TreeMap<KafkaTopicPartition,Long> in method package getRestoredState() : TreeMap<KafkaTopicPartition,Long> in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Return Type KafkaTopicPartitionState<KPH>[] to List<KafkaTopicPartitionState<KPH>> in method private createPartitionStateHolders(partitionsToInitialOffsets Map<KafkaTopicPartition,Long>, timestampWatermarkMode int, watermarksPeriodic SerializedValue<AssignerWithPeriodicWatermarks<T>>, watermarksPunctuated SerializedValue<AssignerWithPunctuatedWatermarks<T>>, userCodeClassLoader ClassLoader) : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Change Attribute Type allPartitions : KafkaTopicPartitionStateWithPeriodicWatermarks<?,?>[] to allPartitions : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher.PeriodicWatermarkEmitter |
From |
To |
|
Change Attribute Type subscribedPartitionStates : KafkaTopicPartitionState<KPH>[] to subscribedPartitionStates : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Change Attribute Type subscribedPartitionStates : KafkaTopicPartitionState<TopicPartition>[] to unassignedPartitionsQueue : ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>> in class org.apache.flink.streaming.connectors.kafka.internal.KafkaConsumerThread |
From |
To |
|
Change Attribute Type restoredState : HashMap<KafkaTopicPartition,Long> to restoredState : TreeMap<KafkaTopicPartition,Long> in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Attribute Type partitionStates : KafkaTopicPartitionState<?>[] to partitionStates : List<KafkaTopicPartitionState<TopicAndPartition>> in class org.apache.flink.streaming.connectors.kafka.internals.PeriodicOffsetCommitter |
From |
To |
|
Rename Parameter subscribedPartitionStates : KafkaTopicPartitionState<TopicPartition>[] to unassignedPartitionsQueue : ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>> in method public KafkaConsumerThread(log Logger, handover Handover, kafkaProperties Properties, unassignedPartitionsQueue ClosableBlockingQueue<KafkaTopicPartitionState<TopicPartition>>, kafkaMetricGroup MetricGroup, consumerCallBridge KafkaConsumerCallBridge, threadName String, pollTimeout long, useMetrics boolean) in class org.apache.flink.streaming.connectors.kafka.internal.KafkaConsumerThread |
From |
To |
|
Rename Parameter assignedPartitionsWithInitialOffsets : Map<KafkaTopicPartition,Long> to seedPartitionsWithInitialOffsets : Map<KafkaTopicPartition,Long> in method protected AbstractFetcher(sourceContext SourceContext<T>, seedPartitionsWithInitialOffsets Map<KafkaTopicPartition,Long>, watermarksPeriodic SerializedValue<AssignerWithPeriodicWatermarks<T>>, watermarksPunctuated SerializedValue<AssignerWithPunctuatedWatermarks<T>>, processingTimeProvider ProcessingTimeService, autoWatermarkInterval long, userCodeClassLoader ClassLoader, useMetrics boolean) in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |
|
Rename Parameter inPartitions : List<KafkaTopicPartition> to partitions : List<KafkaTopicPartition> in method private contains(partitions List<KafkaTopicPartition>, partition int) : boolean in class org.apache.flink.streaming.connectors.kafka.internals.AbstractPartitionDiscovererTest |
From |
To |
|
Rename Parameter assignedPartitionsWithInitialOffsets : Map<KafkaTopicPartition,Long> to seedPartitionsWithInitialOffsets : Map<KafkaTopicPartition,Long> in method public Kafka08Fetcher(sourceContext SourceContext<T>, seedPartitionsWithInitialOffsets Map<KafkaTopicPartition,Long>, watermarksPeriodic SerializedValue<AssignerWithPeriodicWatermarks<T>>, watermarksPunctuated SerializedValue<AssignerWithPunctuatedWatermarks<T>>, runtimeContext StreamingRuntimeContext, deserializer KeyedDeserializationSchema<T>, kafkaProperties Properties, autoCommitInterval long, useMetrics boolean) in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Rename Parameter assignedPartitionsToInitialOffsets : Map<KafkaTopicPartition,Long> to partitionsToInitialOffsets : Map<KafkaTopicPartition,Long> in method private createPartitionStateHolders(partitionsToInitialOffsets Map<KafkaTopicPartition,Long>, timestampWatermarkMode int, watermarksPeriodic SerializedValue<AssignerWithPeriodicWatermarks<T>>, watermarksPunctuated SerializedValue<AssignerWithPunctuatedWatermarks<T>>, userCodeClassLoader ClassLoader) : List<KafkaTopicPartitionState<KPH>> in class org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher |
From |
To |