in flink-connector/flink-connector-gcp-pubsub/src/main/java/com/google/pubsub/flink/internal/source/reader/PubSubSplitReader.java [46:54]
private Multimap<String, PubsubMessage> getMessages() throws Throwable {
ImmutableListMultimap.Builder<String, PubsubMessage> messages = ImmutableListMultimap.builder();
for (Map.Entry<String, NotifyingPullSubscriber> entry : subscribers.entrySet()) {
for (PubsubMessage m : entry.getValue().pullMessage().asSet()) {
messages.put(entry.getKey(), m);
}
}
return messages.build();
}