Github user EAlexRojas commented on a diff in the pull request:

    https://github.com/apache/flink/pull/5991#discussion_r190530647
  
    --- Diff: 
flink-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/KafkaConsumerThread.java
 ---
    @@ -240,7 +249,9 @@ public void run() {
                                                newPartitions = 
unassignedPartitionsQueue.getBatchBlocking();
                                        }
                                        if (newPartitions != null) {
    -                                           
reassignPartitions(newPartitions);
    +                                           
reassignPartitions(newPartitions, new HashSet<>());
    --- End diff --
    
    I just realized this should be actually
    `reassignPartitions(newPartitions, partitionsToBeRemoved);`


---

Reply via email to