|
Change Parameter Type unassignedPartitions : ClosableBlockingQueue<FetchPartition> to unassignedPartitions : ClosableBlockingQueue<KafkaTopicPartitionState<TopicAndPartition>> in method public SimpleConsumerThread(owner Kafka08Fetcher<T>, errorHandler ExceptionProxy, config Properties, broker Node, seedPartitions List<KafkaTopicPartitionState<TopicAndPartition>>, unassignedPartitions ClosableBlockingQueue<KafkaTopicPartitionState<TopicAndPartition>>, deserializer KeyedDeserializationSchema<T>, invalidOffsetBehavior long) in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Parameter Type next : Object to next : Tuple2<Long,String> in method public partition(next Tuple2<Long,String>, serializedKey byte[], serializedValue byte[], numPartitions int) : int in class org.apache.flink.streaming.connectors.kafka.KafkaProducerTestBase.CustomPartitioner |
From |
To |
|
Change Parameter Type inPartitions : List<KafkaTopicPartitionLeader> to inPartitions : List<KafkaTopicPartition> in method private contains(inPartitions List<KafkaTopicPartition>, partition int) : boolean in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Parameter Type partitions : List<T> to allPartitions : List<KafkaTopicPartition> in method protected assignPartitions(allPartitions List<KafkaTopicPartition>, numConsumers int, consumerIndex int) : List<KafkaTopicPartition> in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Parameter Type partitions : List<FetchPartition> to partitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method private getLastOffsetFromKafka(consumer SimpleConsumer, partitions List<KafkaTopicPartitionState<TopicAndPartition>>, whichTime long) : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Parameter Type partitions : List<KafkaTopicPartition> to partitions : KafkaTopicPartitionState<TopicPartition>[] in method public convertKafkaPartitions(partitions KafkaTopicPartitionState<TopicPartition>[]) : List<TopicPartition> in class org.apache.flink.streaming.connectors.kafka.internal.Kafka09Fetcher |
From |
To |
|
Change Parameter Type partitionsList : List<FetchPartition> to partitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method private getTopics(partitions List<KafkaTopicPartitionState<TopicAndPartition>>) : List<String> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Parameter Type partitionsToAssign : List<FetchPartition> to partitionsToAssign : List<KafkaTopicPartitionState<TopicAndPartition>> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Parameter Type owner : LegacyFetcher to owner : Kafka08Fetcher<T> in method public SimpleConsumerThread(owner Kafka08Fetcher<T>, errorHandler ExceptionProxy, config Properties, broker Node, seedPartitions List<KafkaTopicPartitionState<TopicAndPartition>>, unassignedPartitions ClosableBlockingQueue<KafkaTopicPartitionState<TopicAndPartition>>, deserializer KeyedDeserializationSchema<T>, invalidOffsetBehavior long) in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Parameter Type partitions : List<FetchPartition> to partitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method private getMissingOffsetsFromKafka(partitions List<KafkaTopicPartitionState<TopicAndPartition>>) : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Parameter Type seedPartitions : List<FetchPartition> to seedPartitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method public SimpleConsumerThread(owner Kafka08Fetcher<T>, errorHandler ExceptionProxy, config Properties, broker Node, seedPartitions List<KafkaTopicPartitionState<TopicAndPartition>>, unassignedPartitions ClosableBlockingQueue<KafkaTopicPartitionState<TopicAndPartition>>, deserializer KeyedDeserializationSchema<T>, invalidOffsetBehavior long) in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type o2 : long to o2 : Long in method public testOffsetAutocommitTest() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type fp : FetchPartition to part : KafkaTopicPartitionState<TopicAndPartition> in method private getMissingOffsetsFromKafka(partitions List<KafkaTopicPartitionState<TopicAndPartition>>) : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type newPartition : FetchPartition to newPartition : KafkaTopicPartitionState<TopicAndPartition> in method public run() : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type partitionsIterator : Iterator<FetchPartition> to partitionsIterator : Iterator<KafkaTopicPartitionState<TopicAndPartition>> in method public run() : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type newPartitions : List<KafkaTopicPartitionLeader> to newPartitions : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o2 : long to o2 : Long in method public testOffsetInZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type parts2 : List<KafkaTopicPartitionLeader> to parts2 : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type fp : FetchPartition to part : KafkaTopicPartitionState<TopicAndPartition> in method private getLastOffsetFromKafka(consumer SimpleConsumer, partitions List<KafkaTopicPartitionState<TopicAndPartition>>, whichTime long) : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type newPartitions : List<FetchPartition> to newPartitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method public run() : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type parts3 : List<KafkaTopicPartitionLeader> to parts3 : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type allPartitions : Set<KafkaTopicPartitionLeader> to allPartitions : Set<KafkaTopicPartition> in method public testMultiplePartitionsPerConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o2 : long to o2 : Long in method public testKafkaOffsetRetrievalToZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type fp : FetchPartition to fp : KafkaTopicPartitionState<TopicAndPartition> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type unassignedPartitionsIterator : Iterator<FetchPartition> to unassignedPartitionsIterator : Iterator<KafkaTopicPartitionState<TopicAndPartition>> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type parts : List<KafkaTopicPartitionLeader> to parts : List<KafkaTopicPartition> in method public testPartitionsEqualConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o1 : long to o1 : Long in method public testOffsetInZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type partitions : List<FetchPartition> to partitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type leaderToPartitions : Map<Node,List<FetchPartition>> to leaderToPartitions : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type p : KafkaTopicPartitionLeader to p : KafkaTopicPartition in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type partitionsToGetOffsetsFor : List<FetchPartition> to partitionsToGetOffsetsFor : List<KafkaTopicPartitionState<TopicAndPartition>> in method public run() : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type partitionsWithLeaders : Map<Node,List<FetchPartition>> to partitionsWithLeaders : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type parts : List<KafkaTopicPartitionLeader> to parts : List<KafkaTopicPartition> in method public testMultiplePartitionsPerConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o1 : long to o1 : Long in method public testOffsetAutocommitTest() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type partitionsOfLeader : List<FetchPartition> to partitionsOfLeader : List<KafkaTopicPartitionState<TopicAndPartition>> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type parts1 : List<KafkaTopicPartitionLeader> to parts1 : List<KafkaTopicPartition> in method public testAssignEmptyPartitions() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type p : KafkaTopicPartitionLeader to p : KafkaTopicPartition in method public testMultiplePartitionsPerConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type allNewPartitions : Set<KafkaTopicPartitionLeader> to allNewPartitions : Set<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o1 : long to o1 : Long in method public testKafkaOffsetRetrievalToZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type partitionsToGetOffsetsFor : List<FetchPartition> to partitionsToGetOffsetsFor : List<KafkaTopicPartitionState<TopicAndPartition>> in method private getMissingOffsetsFromKafka(partitions List<KafkaTopicPartitionState<TopicAndPartition>>) : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type unassignedPartition : FetchPartition to unassignedPartition : KafkaTopicPartitionState<TopicAndPartition> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type offset : long to offset : Long in method public getOffsets(partitions List<KafkaTopicPartition>) : Map<KafkaTopicPartition,Long> in class org.apache.flink.streaming.connectors.kafka.internals.ZookeeperOffsetHandler |
From |
To |
|
Change Variable Type partitionsWithLeader : Map.Entry<Node,List<FetchPartition>> to partitionsWithLeader : Map.Entry<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type ktp : KafkaTopicPartitionLeader to ktp : KafkaTopicPartition in method private contains(inPartitions List<KafkaTopicPartition>, partition int) : boolean in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type ep : List<KafkaTopicPartitionLeader> to ep : List<KafkaTopicPartition> in method public testAssignEmptyPartitions() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type p : KafkaTopicPartitionLeader to p : KafkaTopicPartition in method public testPartitionsFewerThanConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type o3 : long to o3 : Long in method public testOffsetAutocommitTest() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type partitionsToSub : List<T> to thisSubtaskPartitions : List<KafkaTopicPartition> in method protected assignPartitions(allPartitions List<KafkaTopicPartition>, numConsumers int, consumerIndex int) : List<KafkaTopicPartition> in class org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase |
From |
To |
|
Change Variable Type initialPartitions : List<KafkaTopicPartitionLeader> to initialPartitions : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type allInitialPartitions : Set<KafkaTopicPartitionLeader> to allInitialPartitions : Set<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type fp : FetchPartition to fp : KafkaTopicPartitionState<TopicAndPartition> in method public run() : void in class org.apache.flink.streaming.connectors.kafka.internals.SimpleConsumerThread |
From |
To |
|
Change Variable Type parts : List<KafkaTopicPartitionLeader> to parts : List<KafkaTopicPartition> in method public testPartitionsFewerThanConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type allPartitions : Set<KafkaTopicPartitionLeader> to allPartitions : Set<KafkaTopicPartition> in method public testPartitionsFewerThanConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type unassignedPartitions : List<FetchPartition> to unassignedPartitions : List<KafkaTopicPartitionState<TopicAndPartition>> in method private findLeaderForPartitions(partitionsToAssign List<KafkaTopicPartitionState<TopicAndPartition>>, kafkaProperties Properties) : Map<Node,List<KafkaTopicPartitionState<TopicAndPartition>>> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type o3 : long to o3 : Long in method public testKafkaOffsetRetrievalToZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type partitionsToAssign : List<FetchPartition> to partitionsToAssign : List<KafkaTopicPartitionState<TopicAndPartition>> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type fp : FetchPartition to fp : KafkaTopicPartitionState<TopicAndPartition> in method private getTopics(partitions List<KafkaTopicPartitionState<TopicAndPartition>>) : List<String> in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type o3 : long to o3 : Long in method public testOffsetInZookeeper() : void in class org.apache.flink.streaming.connectors.kafka.Kafka08ITCase |
From |
To |
|
Change Variable Type newPartitionsQueue : ClosableBlockingQueue<FetchPartition> to newPartitionsQueue : ClosableBlockingQueue<KafkaTopicPartitionState<TopicAndPartition>> in method public runFetchLoop() : void in class org.apache.flink.streaming.connectors.kafka.internals.Kafka08Fetcher |
From |
To |
|
Change Variable Type partitions : List<KafkaTopicPartitionLeader> to partitions : List<KafkaTopicPartition> in method public testMultiplePartitionsPerConsumers() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type parts3new : List<KafkaTopicPartitionLeader> to parts3new : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type parts2new : List<KafkaTopicPartitionLeader> to parts2new : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type parts2 : List<KafkaTopicPartitionLeader> to parts2 : List<KafkaTopicPartition> in method public testAssignEmptyPartitions() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type parts1new : List<KafkaTopicPartitionLeader> to parts1new : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |
|
Change Variable Type parts1 : List<KafkaTopicPartitionLeader> to parts1 : List<KafkaTopicPartition> in method public testGrowingPartitionsRemainsStable() : void in class org.apache.flink.streaming.connectors.kafka.KafkaConsumerPartitionAssignmentTest |
From |
To |