src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java [361:368]:
- - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - -
            Map<MessageQueue, Long> result = new ConcurrentHashMap<>();
            CompletableFuture.allOf(
                            messageQueues.stream()
                                    .map(
                                            messageQueue ->
                                                    CompletableFuture.supplyAsync(
                                                                    () ->
                                                                            innerConsumer
- - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - -



src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java [419:426]:
- - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - -
            Map<MessageQueue, Long> result = new ConcurrentHashMap<>();
            CompletableFuture.allOf(
                            messageQueues.stream()
                                    .map(
                                            messageQueue ->
                                                    CompletableFuture.supplyAsync(
                                                                    () ->
                                                                            innerConsumer
- - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - -



