This is an automated email from the ASF dual-hosted git repository.

Arsnael pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git


The following commit(s) were added to refs/heads/master by this push:
     new dc7900624a [FIX] PG: Generalise streaming with paging
dc7900624a is described below

commit dc7900624aaddd066f34c9e2ba76d6e38d259769
Author: Benoit TELLIER <[email protected]>
AuthorDate: Fri Sep 25 15:14:16 2026 +0200

    [FIX] PG: Generalise streaming with paging
---
 .../backends/postgres/utils/PostgresExecutor.java  |  38 +++
 .../james/events/PostgresEventDeadLetters.java     |  11 +-
 .../PostgresDeletedMessageMetadataVault.java       |  11 +-
 .../postgres/mail/dao/PostgresAttachmentDAO.java   |   7 +-
 .../postgres/mail/dao/PostgresMailboxDAO.java      |   4 +-
 .../mail/dao/PostgresMailboxMessageDAO.java        | 262 ++++++---------------
 .../postgres/mail/dao/PostgresMessageDAO.java      |   7 +-
 .../james/blob/postgres/PostgresBlobStoreDAO.java  |  25 +-
 .../jmap/postgres/upload/PostgresUploadDAO.java    |   7 +-
 .../postgres/PostgresMailRepositoryContentDAO.java |  13 +-
 .../postgres/PostgresRecipientRewriteTableDAO.java |   7 +-
 .../james/user/postgres/PostgresUsersDAO.java      |  19 +-
 12 files changed, 171 insertions(+), 240 deletions(-)

diff --git 
a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
 
b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
index cbe00f8023..60b1516d6c 100644
--- 
a/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
+++ 
b/backends-common/postgres/src/main/java/org/apache/james/backends/postgres/utils/PostgresExecutor.java
@@ -23,8 +23,10 @@ import static org.jooq.impl.DSL.exists;
 import static org.jooq.impl.DSL.field;
 
 import java.time.Duration;
+import java.util.List;
 import java.util.Optional;
 import java.util.concurrent.TimeoutException;
+import java.util.function.BiFunction;
 import java.util.function.Function;
 import java.util.function.Predicate;
 
@@ -39,6 +41,7 @@ import org.jooq.Record;
 import org.jooq.Record1;
 import org.jooq.SQLDialect;
 import org.jooq.SelectConditionStep;
+import org.jooq.SelectLimitStep;
 import org.jooq.conf.Settings;
 import org.jooq.conf.StatementType;
 import org.jooq.impl.DSL;
@@ -163,6 +166,41 @@ public class PostgresExecutor {
                 jamesPostgresConnectionFactory::closeConnection)));
     }
 
+    /**
+     * Streams a potentially large result set using keyset pagination.
+     * <p>
+     * Each page is a short query, fully read before being emitted: the 
connection is released between pages, no cursor nor
+     * transaction is held open, and the next page is only fetched upon 
downstream demand. This makes it safe for slow
+     * consumers (throttled tasks, re-indexing...) while keeping memory 
bounded to a couple of pages.
+     * <p>
+     * Streaming a single query instead would hold a connection (and a 
Postgres snapshot) for the whole duration of the
+     * consumption, and would trip the reactive timeout as soon as the 
consumer pauses for longer than it.
+     *
+     * @param pageQuery builds the query for a page given the last record of 
the previous page (empty for the first page).
+     *                  It MUST order results by a unique key and filter 
records strictly after the given record for that key.
+     *                  The limit is applied by this method.
+     */
+    public Flux<Record> executeRowsPaginated(BiFunction<DSLContext, 
Optional<Record>, SelectLimitStep<? extends Record>> pageQuery) {
+        return executeRowsPaginated(pageQuery, PostgresUtils.QUERY_BATCH_SIZE);
+    }
+
+    public Flux<Record> executeRowsPaginated(BiFunction<DSLContext, 
Optional<Record>, SelectLimitStep<? extends Record>> pageQuery, int pageSize) {
+        return Flux.defer(() -> executePage(pageQuery, Optional.empty(), 
pageSize))
+            .expand(page -> {
+                if (page.size() < pageSize) {
+                    return Mono.empty();
+                }
+                return executePage(pageQuery, Optional.of(page.getLast()), 
pageSize);
+            })
+            // prefetch 1 page: do not read ahead more than what is needed
+            .concatMapIterable(Function.identity(), 1);
+    }
+
+    private Mono<List<Record>> executePage(BiFunction<DSLContext, 
Optional<Record>, SelectLimitStep<? extends Record>> pageQuery, 
Optional<Record> lastRecord, int pageSize) {
+        return executeRows(dslContext -> Flux.from(pageQuery.apply(dslContext, 
lastRecord).limit(pageSize)))
+            .collectList();
+    }
+
     public Flux<Record> executeDeleteAndReturnList(Function<DSLContext, 
DeleteResultStep<Record>> queryFunction) {
         return 
Flux.from(metricFactory.decoratePublisherWithTimerMetric("postgres-execution",
             Flux.usingWhen(getConnection(domain),
diff --git 
a/event-bus/postgres/src/main/java/org/apache/james/events/PostgresEventDeadLetters.java
 
b/event-bus/postgres/src/main/java/org/apache/james/events/PostgresEventDeadLetters.java
index de967e62b9..ece3e1d337 100644
--- 
a/event-bus/postgres/src/main/java/org/apache/james/events/PostgresEventDeadLetters.java
+++ 
b/event-bus/postgres/src/main/java/org/apache/james/events/PostgresEventDeadLetters.java
@@ -28,6 +28,7 @@ import jakarta.inject.Inject;
 
 import org.apache.james.backends.postgres.utils.PostgresExecutor;
 import org.jooq.Record;
+import org.jooq.impl.DSL;
 
 import com.github.fge.lambdas.Throwing;
 import com.google.common.base.Preconditions;
@@ -99,10 +100,12 @@ public class PostgresEventDeadLetters implements 
EventDeadLetters {
     public Flux<InsertionId> failedIds(Group registeredGroup) {
         Preconditions.checkArgument(registeredGroup != null, 
REGISTERED_GROUP_CANNOT_BE_NULL);
 
-        return postgresExecutor.executeRows(dslContext -> Flux.from(dslContext
-            .select(INSERTION_ID)
-            .from(TABLE_NAME)
-            .where(GROUP.eq(registeredGroup.asString()))))
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext
+                .select(INSERTION_ID)
+                .from(TABLE_NAME)
+                .where(GROUP.eq(registeredGroup.asString()))
+                .and(lastRecord.map(record -> 
INSERTION_ID.greaterThan(record.get(INSERTION_ID))).orElseGet(DSL::noCondition))
+                .orderBy(INSERTION_ID))
             .map(record -> InsertionId.of(record.get(INSERTION_ID)));
     }
 
diff --git 
a/mailbox/plugin/deleted-messages-vault-postgres/src/main/java/org/apache/james/vault/metadata/PostgresDeletedMessageMetadataVault.java
 
b/mailbox/plugin/deleted-messages-vault-postgres/src/main/java/org/apache/james/vault/metadata/PostgresDeletedMessageMetadataVault.java
index 38a5bc9e6c..e25bd8e18c 100644
--- 
a/mailbox/plugin/deleted-messages-vault-postgres/src/main/java/org/apache/james/vault/metadata/PostgresDeletedMessageMetadataVault.java
+++ 
b/mailbox/plugin/deleted-messages-vault-postgres/src/main/java/org/apache/james/vault/metadata/PostgresDeletedMessageMetadataVault.java
@@ -38,6 +38,7 @@ import org.apache.james.blob.api.BucketName;
 import org.apache.james.core.Username;
 import org.apache.james.mailbox.model.MessageId;
 import org.jooq.Record;
+import org.jooq.impl.DSL;
 import org.reactivestreams.Publisher;
 
 import reactor.core.publisher.Flux;
@@ -98,10 +99,12 @@ public class PostgresDeletedMessageMetadataVault implements 
DeletedMessageMetada
 
     @Override
     public Publisher<DeletedMessageWithStorageInformation> 
listMessages(BucketName bucketName, Username username) {
-        return postgresExecutor.executeRows(context -> 
Flux.from(context.select(METADATA)
-            .from(TABLE_NAME)
-            .where(BUCKET_NAME.eq(bucketName.asString()),
-                OWNER.eq(username.asString()))))
+        return postgresExecutor.executeRowsPaginated((context, lastRecord) -> 
context.select(MESSAGE_ID, METADATA)
+                .from(TABLE_NAME)
+                .where(BUCKET_NAME.eq(bucketName.asString()),
+                    OWNER.eq(username.asString()))
+                .and(lastRecord.map(record -> 
MESSAGE_ID.greaterThan(record.get(MESSAGE_ID))).orElseGet(DSL::noCondition))
+                .orderBy(MESSAGE_ID))
             .map(record -> 
metadataSerializer.deserialize(record.get(METADATA).data()))
             .handle(publishIfPresent());
     }
diff --git 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresAttachmentDAO.java
 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresAttachmentDAO.java
index 8439ec029d..96fb9bfcab 100644
--- 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresAttachmentDAO.java
+++ 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresAttachmentDAO.java
@@ -34,6 +34,7 @@ import org.apache.james.mailbox.model.AttachmentMetadata;
 import org.apache.james.mailbox.model.StringBackedAttachmentId;
 import org.apache.james.mailbox.postgres.PostgresMessageId;
 import 
org.apache.james.mailbox.postgres.mail.PostgresAttachmentDataDefinition.PostgresAttachmentTable;
+import org.jooq.impl.DSL;
 
 import com.google.common.collect.ImmutableList;
 
@@ -120,8 +121,10 @@ public class PostgresAttachmentDAO {
     }
 
     public Flux<BlobId> listBlobs() {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(PostgresAttachmentTable.BLOB_ID)
-                .from(PostgresAttachmentTable.TABLE_NAME)))
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select(PostgresAttachmentTable.ID, 
PostgresAttachmentTable.BLOB_ID)
+                .from(PostgresAttachmentTable.TABLE_NAME)
+                .where(lastRecord.map(record -> 
PostgresAttachmentTable.ID.greaterThan(record.get(PostgresAttachmentTable.ID))).orElseGet(DSL::noCondition))
+                .orderBy(PostgresAttachmentTable.ID))
             .map(row -> 
blobIdFactory.parse(row.get(PostgresAttachmentTable.BLOB_ID)));
     }
 }
\ No newline at end of file
diff --git 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxDAO.java
 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxDAO.java
index 2e56b6373e..8bbc915339 100644
--- 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxDAO.java
+++ 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxDAO.java
@@ -253,7 +253,9 @@ public class PostgresMailboxDAO {
     }
 
     public Flux<PostgresMailbox> getAll() {
-        return postgresExecutor.executeRows(dsl -> 
Flux.from(dsl.selectFrom(TABLE_NAME)))
+        return postgresExecutor.executeRowsPaginated((dsl, lastRecord) -> 
dsl.selectFrom(TABLE_NAME)
+                .where(lastRecord.map(record -> 
MAILBOX_ID.greaterThan(record.get(MAILBOX_ID))).orElseGet(DSL::noCondition))
+                .orderBy(MAILBOX_ID))
             .map(RECORD_TO_POSTGRES_MAILBOX_FUNCTION);
     }
 
diff --git 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxMessageDAO.java
 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxMessageDAO.java
index 22197fe710..27e4fb258f 100644
--- 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxMessageDAO.java
+++ 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMailboxMessageDAO.java
@@ -64,7 +64,6 @@ import jakarta.mail.Flags;
 
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.james.backends.postgres.utils.PostgresExecutor;
-import org.apache.james.backends.postgres.utils.PostgresUtils;
 import org.apache.james.core.Domain;
 import org.apache.james.mailbox.MessageUid;
 import org.apache.james.mailbox.ModSeq;
@@ -118,8 +117,6 @@ public class PostgresMailboxMessageDAO {
 
     public static final SortField<Long> DEFAULT_SORT_ORDER_BY = 
MESSAGE_UID.asc();
 
-    private static final int QUERY_BATCH_SIZE = PostgresUtils.QUERY_BATCH_SIZE;
-
     private final PostgresExecutor postgresExecutor;
 
     public PostgresMailboxMessageDAO(PostgresExecutor postgresExecutor) {
@@ -137,32 +134,20 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<MessageUid> listUnseen(PostgresMailboxId mailboxId) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq((mailboxId.asUuid())))
-                .and(IS_SEEN.eq(false))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_SEEN.eq(false)));
     }
 
     public Flux<MessageUid> listUnseen(PostgresMailboxId mailboxId, 
MessageRange range) {
         return switch (range.getType()) {
             case ALL -> listUnseen(mailboxId);
-            case FROM -> postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                    .from(TABLE_NAME)
-                    .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                    .and(IS_SEEN.eq(false))
-                    
.and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
-                    .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-                .map(RECORD_TO_MESSAGE_UID_FUNCTION);
-            case RANGE -> postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                    .from(TABLE_NAME)
-                    .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                    .and(IS_SEEN.eq(false))
-                    
.and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
-                    .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong()))
-                    .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-                .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+            case FROM -> listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+                .and(IS_SEEN.eq(false))
+                .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong())));
+            case RANGE -> listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+                .and(IS_SEEN.eq(false))
+                .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
+                .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong())));
             case ONE -> postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
                     .from(TABLE_NAME)
                     .where(MAILBOX_ID.eq(mailboxId.asUuid()))
@@ -177,20 +162,12 @@ public class PostgresMailboxMessageDAO {
         if (!StoreMessageManager.HANDLE_RECENT) {
             return Flux.empty();
         }
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq((mailboxId.asUuid())))
-                .and(IS_RECENT.eq(true))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_RECENT.eq(true)));
     }
 
     public Flux<MessageUid> listAllMessageUid(PostgresMailboxId mailboxId) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq((mailboxId.asUuid())))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid()));
     }
 
     public Flux<MessageUid> listUids(PostgresMailboxId mailboxId, MessageRange 
range) {
@@ -201,15 +178,25 @@ public class PostgresMailboxMessageDAO {
     }
 
     private Flux<MessageUid> doListUids(PostgresMailboxId mailboxId, 
MessageRange range) {
-        return  postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
+            .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong())));
+    }
+
+    private Flux<MessageUid> listUidsPaginated(Condition condition) {
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select(MESSAGE_UID)
                 .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
-                .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong()))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
+                .where(condition)
+                .and(afterUid(lastRecord))
+                .orderBy(DEFAULT_SORT_ORDER_BY))
             .map(RECORD_TO_MESSAGE_UID_FUNCTION);
     }
 
+    private static Condition afterUid(Optional<Record> lastRecord) {
+        return lastRecord.map(record -> 
MESSAGE_UID.greaterThan(record.get(MESSAGE_UID)))
+            .orElseGet(DSL::noCondition);
+    }
+
     public Mono<MessageMetaData> 
deleteByMailboxIdAndMessageUid(PostgresMailboxId mailboxId, MessageUid 
messageUid) {
         return postgresExecutor.executeRow(dslContext -> 
Mono.from(dslContext.deleteFrom(TABLE_NAME)
                 .where(MAILBOX_ID.eq(mailboxId.asUuid()))
@@ -277,66 +264,30 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<Pair<SimpleMailboxMessage.Builder, Record>> 
findMessagesByMailboxId(PostgresMailboxId mailboxId, Limit limit, 
MessageMapper.FetchType fetchType) {
-        if (limit.isUnlimited()) {
-            return Flux.defer(() -> findMessagesByMailboxIdBatch(mailboxId, 
fetchType, Optional.empty(), QUERY_BATCH_SIZE))
-                .expand(messages -> {
-                    if (messages.isEmpty() || messages.size() < 
QUERY_BATCH_SIZE) {
-                        return Mono.empty();
-                    }
-                    return findMessagesByMailboxIdBatch(mailboxId, fetchType, 
Optional.of(messages.getLast().getRight().get(MESSAGE_UID)), QUERY_BATCH_SIZE);
-                })
-                .flatMapIterable(Function.identity());
-        } else {
-            return findMessagesByMailboxIdBatch(mailboxId, fetchType, 
Optional.empty(), limit.getLimit().get())
-                .flatMapIterable(Function.identity());
-        }
-    }
-
-    private Mono<List<Pair<SimpleMailboxMessage.Builder, Record>>> 
findMessagesByMailboxIdBatch(PostgresMailboxId mailboxId, 
MessageMapper.FetchType fetchType,
-                                                                               
                 Optional<Long> messageUidFrom, int batchSize) {
-        PostgresMailboxMessageFetchStrategy fetchStrategy = 
FETCH_TYPE_TO_FETCH_STRATEGY.apply(fetchType);
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(fetchStrategy.fetchFields())
-                .from(MESSAGES_JOIN_MAILBOX_MESSAGES_CONDITION_STEP)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                
.and(messageUidFrom.map(MESSAGE_UID::greaterThan).orElseGet(DSL::noCondition))
-                .orderBy(MESSAGE_UID.asc())
-                .limit(batchSize)))
-            .map(record -> 
Pair.of(fetchStrategy.toMessageBuilder().apply(record), record))
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+        return findMessagesPaginated(MAILBOX_ID.eq(mailboxId.asUuid()), limit, 
fetchType);
     }
 
     public Flux<Pair<SimpleMailboxMessage.Builder, Record>> 
findMessagesByMailboxIdAndBetweenUIDs(PostgresMailboxId mailboxId, MessageUid 
from, MessageUid to, Limit limit, FetchType fetchType) {
-        if (limit.isUnlimited()) {
-            return Flux.defer(() -> 
findMessagesByMailboxIdAndBetweenUIDsBatch(mailboxId, 
MESSAGE_UID.greaterOrEqual(from.asLong()), to, fetchType, QUERY_BATCH_SIZE))
-                .expand(messages -> {
-                    if (messages.isEmpty() || messages.size() < 
QUERY_BATCH_SIZE) {
-                        return Mono.empty();
-                    }
-                    MessageUid messageUidFrom = 
MessageUid.of(messages.getLast().getRight().get(MESSAGE_UID));
-                    return 
findMessagesByMailboxIdAndBetweenUIDsBatch(mailboxId, 
MESSAGE_UID.greaterThan(messageUidFrom.asLong()), to, fetchType, 
QUERY_BATCH_SIZE);
-                })
-                .flatMapIterable(Function.identity());
-        } else {
-            return findMessagesByMailboxIdAndBetweenUIDsBatch(mailboxId, 
MESSAGE_UID.greaterOrEqual(from.asLong()), to, fetchType, 
limit.getLimit().get())
-                .flatMapIterable(Function.identity());
-        }
+        return findMessagesPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+                .and(MESSAGE_UID.greaterOrEqual(from.asLong()))
+                .and(MESSAGE_UID.lessOrEqual(to.asLong())),
+            limit, fetchType);
     }
 
-    private Mono<List<Pair<SimpleMailboxMessage.Builder, Record>>> 
findMessagesByMailboxIdAndBetweenUIDsBatch(PostgresMailboxId mailboxId, 
Condition messageUidFromCondition,
-                                                                               
                               MessageUid to,
-                                                                               
                               FetchType fetchType, int batchSize) {
+    private Flux<Pair<SimpleMailboxMessage.Builder, Record>> 
findMessagesPaginated(Condition condition, Limit limit, FetchType fetchType) {
         PostgresMailboxMessageFetchStrategy fetchStrategy = 
FETCH_TYPE_TO_FETCH_STRATEGY.apply(fetchType);
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(fetchStrategy.fetchFields())
+        Flux<Record> records = limit.getLimit()
+            .map(limitValue -> postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(fetchStrategy.fetchFields())
                 .from(MESSAGES_JOIN_MAILBOX_MESSAGES_CONDITION_STEP)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(messageUidFromCondition)
-                .and(MESSAGE_UID.lessOrEqual(to.asLong()))
-                .orderBy(MESSAGE_UID.asc())
-                .limit(batchSize)))
-            .map(record -> 
Pair.of(fetchStrategy.toMessageBuilder().apply(record), record))
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+                .where(condition)
+                .orderBy(DEFAULT_SORT_ORDER_BY)
+                .limit(limitValue)), EAGER_FETCH))
+            .orElseGet(() -> 
postgresExecutor.executeRowsPaginated((dslContext, lastRecord) -> 
dslContext.select(fetchStrategy.fetchFields())
+                .from(MESSAGES_JOIN_MAILBOX_MESSAGES_CONDITION_STEP)
+                .where(condition)
+                .and(afterUid(lastRecord))
+                .orderBy(DEFAULT_SORT_ORDER_BY)));
+        return records.map(record -> 
Pair.of(fetchStrategy.toMessageBuilder().apply(record), record));
     }
 
     public Mono<Pair<SimpleMailboxMessage.Builder, Record>> 
findMessageByMailboxIdAndUid(PostgresMailboxId mailboxId, MessageUid uid, 
FetchType fetchType) {
@@ -349,35 +300,9 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<Pair<SimpleMailboxMessage.Builder, Record>> 
findMessagesByMailboxIdAndAfterUID(PostgresMailboxId mailboxId, MessageUid 
from, Limit limit, FetchType fetchType) {
-        if (limit.isUnlimited()) {
-            return Flux.defer(() -> 
findMessagesByMailboxIdAndAfterUIDBatch(mailboxId, 
MESSAGE_UID.greaterOrEqual(from.asLong()), fetchType, QUERY_BATCH_SIZE))
-                .expand(messages -> {
-                    if (messages.isEmpty() || messages.size() < 
QUERY_BATCH_SIZE) {
-                        return Mono.empty();
-                    }
-                    MessageUid messageUidFrom = 
MessageUid.of(messages.getLast().getRight().get(MESSAGE_UID));
-                    return findMessagesByMailboxIdAndAfterUIDBatch(mailboxId, 
MESSAGE_UID.greaterThan(messageUidFrom.asLong()), fetchType, QUERY_BATCH_SIZE);
-                })
-                .flatMapIterable(Function.identity());
-        } else {
-            return findMessagesByMailboxIdAndAfterUIDBatch(mailboxId, 
MESSAGE_UID.greaterOrEqual(from.asLong()), fetchType, limit.getLimit().get())
-                .flatMapIterable(Function.identity());
-        }
-    }
-
-    private Mono<List<Pair<SimpleMailboxMessage.Builder, Record>>> 
findMessagesByMailboxIdAndAfterUIDBatch(PostgresMailboxId mailboxId,
-                                                                               
                            Condition messageUidFromCondition,
-                                                                               
                            FetchType fetchType, int batchSize) {
-        PostgresMailboxMessageFetchStrategy fetchStrategy = 
FETCH_TYPE_TO_FETCH_STRATEGY.apply(fetchType);
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(fetchStrategy.fetchFields())
-                .from(MESSAGES_JOIN_MAILBOX_MESSAGES_CONDITION_STEP)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(messageUidFromCondition)
-                .orderBy(MESSAGE_UID.asc())
-                .limit(batchSize)))
-            .map(record -> 
Pair.of(fetchStrategy.toMessageBuilder().apply(record), record))
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+        return findMessagesPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+                .and(MESSAGE_UID.greaterOrEqual(from.asLong())),
+            limit, fetchType);
     }
 
     public Flux<SimpleMailboxMessage.Builder> 
findMessagesByMailboxIdAndUIDs(PostgresMailboxId mailboxId, List<MessageUid> 
uids) {
@@ -402,33 +327,21 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<MessageUid> findDeletedMessagesByMailboxId(PostgresMailboxId 
mailboxId) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(IS_DELETED.eq(true))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_DELETED.eq(true)));
     }
 
     public Flux<MessageUid> 
findDeletedMessagesByMailboxIdAndBetweenUIDs(PostgresMailboxId mailboxId, 
MessageUid from, MessageUid to) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(IS_DELETED.eq(true))
-                .and(MESSAGE_UID.greaterOrEqual(from.asLong()))
-                .and(MESSAGE_UID.lessOrEqual(to.asLong()))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_DELETED.eq(true))
+            .and(MESSAGE_UID.greaterOrEqual(from.asLong()))
+            .and(MESSAGE_UID.lessOrEqual(to.asLong())));
     }
 
     public Flux<MessageUid> 
findDeletedMessagesByMailboxIdAndAfterUID(PostgresMailboxId mailboxId, 
MessageUid from) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(IS_DELETED.eq(true))
-                .and(MESSAGE_UID.greaterOrEqual(from.asLong()))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_DELETED.eq(true))
+            .and(MESSAGE_UID.greaterOrEqual(from.asLong())));
     }
 
     public Mono<MessageUid> 
findDeletedMessageByMailboxIdAndUid(PostgresMailboxId mailboxId, MessageUid 
uid) {
@@ -441,14 +354,10 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<MessageUid> listNotDeletedUids(PostgresMailboxId mailboxId, 
MessageRange range) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(MESSAGE_UID, IS_DELETED)
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
-                .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong()))
-                .and(IS_DELETED.eq(false))
-                .orderBy(DEFAULT_SORT_ORDER_BY)), EAGER_FETCH)
-            .map(RECORD_TO_MESSAGE_UID_FUNCTION);
+        return listUidsPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
+            .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong()))
+            .and(IS_DELETED.eq(false)));
     }
 
     public Mono<Boolean> existsByMessageId(PostgresMessageId messageId) {
@@ -458,52 +367,23 @@ public class PostgresMailboxMessageDAO {
     }
 
     public Flux<ComposedMessageIdWithMetaData> 
findMessagesMetadata(PostgresMailboxId mailboxId, MessageRange range) {
-        return Flux.defer(() -> findMessagesMetadataBatch(mailboxId, 
range.getUidTo(), MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()), 
QUERY_BATCH_SIZE))
-            .expand(messages -> {
-                if (messages.isEmpty() || messages.size() < QUERY_BATCH_SIZE) {
-                    return Mono.empty();
-                }
-                MessageUid messageUidFrom = 
messages.getLast().getComposedMessageId().getUid();
-                return findMessagesMetadataBatch(mailboxId, range.getUidTo(), 
MESSAGE_UID.greaterThan(messageUidFrom.asLong()), QUERY_BATCH_SIZE);
-            })
-            .flatMapIterable(Function.identity());
-    }
-
-    private Mono<List<ComposedMessageIdWithMetaData>> 
findMessagesMetadataBatch(PostgresMailboxId mailboxId, MessageUid messageUidTo, 
Condition messageUidFromCondition, int batchSize) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select()
-                .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(messageUidFromCondition)
-                .and(MESSAGE_UID.lessOrEqual(messageUidTo.asLong()))
-                .orderBy(MESSAGE_UID.asc())
-                .limit(batchSize)))
-            .map(RECORD_TO_COMPOSED_MESSAGE_ID_WITH_META_DATA_FUNCTION)
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+        return findMessagesMetadataPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(MESSAGE_UID.greaterOrEqual(range.getUidFrom().asLong()))
+            .and(MESSAGE_UID.lessOrEqual(range.getUidTo().asLong())));
     }
 
     public Flux<ComposedMessageIdWithMetaData> 
findAllRecentMessageMetadata(PostgresMailboxId mailboxId) {
-        return Flux.defer(() -> findAllRecentMessageMetadataBatch(mailboxId, 
Optional.empty(), QUERY_BATCH_SIZE))
-            .expand(messages -> {
-                if (messages.isEmpty() || messages.size() < QUERY_BATCH_SIZE) {
-                    return Mono.empty();
-                }
-                return findAllRecentMessageMetadataBatch(mailboxId, 
Optional.of(messages.getLast().getComposedMessageId().getUid()), 
QUERY_BATCH_SIZE);
-            })
-            .flatMapIterable(Function.identity());
+        return findMessagesMetadataPaginated(MAILBOX_ID.eq(mailboxId.asUuid())
+            .and(IS_RECENT.eq(true)));
     }
 
-    private Mono<List<ComposedMessageIdWithMetaData>> 
findAllRecentMessageMetadataBatch(PostgresMailboxId mailboxId, 
Optional<MessageUid> messageUidFrom, int batchSize) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select()
+    private Flux<ComposedMessageIdWithMetaData> 
findMessagesMetadataPaginated(Condition condition) {
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select()
                 .from(TABLE_NAME)
-                .where(MAILBOX_ID.eq(mailboxId.asUuid()))
-                .and(IS_RECENT.eq(true))
-                .and(messageUidFrom.map(messageUid -> 
MESSAGE_UID.greaterThan(messageUid.asLong())).orElseGet(DSL::noCondition))
-                .orderBy(MESSAGE_UID.asc())
-                .limit(batchSize)))
-            .map(RECORD_TO_COMPOSED_MESSAGE_ID_WITH_META_DATA_FUNCTION)
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+                .where(condition)
+                .and(afterUid(lastRecord))
+                .orderBy(DEFAULT_SORT_ORDER_BY))
+            .map(RECORD_TO_COMPOSED_MESSAGE_ID_WITH_META_DATA_FUNCTION);
     }
 
     public Mono<Flags> replaceFlags(PostgresMailboxId mailboxId, MessageUid 
uid, Flags newFlags, ModSeq newModSeq) {
diff --git 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMessageDAO.java
 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMessageDAO.java
index 37ac7e5d49..5d7d2d4373 100644
--- 
a/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMessageDAO.java
+++ 
b/mailbox/postgres/src/main/java/org/apache/james/mailbox/postgres/mail/dao/PostgresMessageDAO.java
@@ -49,6 +49,7 @@ import 
org.apache.james.mailbox.postgres.mail.PostgresMessageDataDefinition;
 import org.apache.james.mailbox.postgres.mail.dto.AttachmentsDTO;
 import org.apache.james.mailbox.store.mail.model.MailboxMessage;
 import org.jooq.Record;
+import org.jooq.impl.DSL;
 
 import reactor.core.publisher.Flux;
 import reactor.core.publisher.Mono;
@@ -126,8 +127,10 @@ public class PostgresMessageDAO {
     }
 
     public Flux<BlobId> listBlobs() {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(BODY_BLOB_ID)
-            .from(TABLE_NAME)))
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select(MESSAGE_ID, BODY_BLOB_ID)
+                .from(TABLE_NAME)
+                .where(lastRecord.map(record -> 
MESSAGE_ID.greaterThan(record.get(MESSAGE_ID))).orElseGet(DSL::noCondition))
+                .orderBy(MESSAGE_ID))
             .map(record -> blobIdFactory.parse(record.get(BODY_BLOB_ID)));
     }
 
diff --git 
a/server/blob/blob-postgres/src/main/java/org/apache/james/blob/postgres/PostgresBlobStoreDAO.java
 
b/server/blob/blob-postgres/src/main/java/org/apache/james/blob/postgres/PostgresBlobStoreDAO.java
index fc20507a36..5f15e25377 100644
--- 
a/server/blob/blob-postgres/src/main/java/org/apache/james/blob/postgres/PostgresBlobStoreDAO.java
+++ 
b/server/blob/blob-postgres/src/main/java/org/apache/james/blob/postgres/PostgresBlobStoreDAO.java
@@ -29,16 +29,13 @@ import static 
org.apache.james.blob.postgres.PostgresBlobStorageDataDefinition.P
 import java.io.IOException;
 import java.io.InputStream;
 import java.util.Collection;
-import java.util.List;
 import java.util.Map;
 import java.util.Optional;
-import java.util.function.Function;
 
 import jakarta.inject.Inject;
 
 import org.apache.commons.io.IOUtils;
 import org.apache.james.backends.postgres.utils.PostgresExecutor;
-import org.apache.james.backends.postgres.utils.PostgresUtils;
 import org.apache.james.blob.api.BlobId;
 import org.apache.james.blob.api.BlobStoreDAO;
 import org.apache.james.blob.api.BucketName;
@@ -168,26 +165,12 @@ public class PostgresBlobStoreDAO implements BlobStoreDAO 
{
 
     @Override
     public Flux<BlobId> listBlobs(BucketName bucketName) {
-        return Flux.defer(() -> listBlobsBatch(bucketName, Optional.empty(), 
PostgresUtils.QUERY_BATCH_SIZE))
-            .expand(blobIds -> {
-                if (blobIds.isEmpty() || blobIds.size() < 
PostgresUtils.QUERY_BATCH_SIZE) {
-                    return Mono.empty();
-                }
-                return listBlobsBatch(bucketName, 
Optional.of(blobIds.getLast()), PostgresUtils.QUERY_BATCH_SIZE);
-            })
-            .flatMapIterable(Function.identity());
-    }
-
-    private Mono<List<BlobId>> listBlobsBatch(BucketName bucketName, 
Optional<BlobId> blobIdFrom, int batchSize) {
-        return postgresExecutor.executeRows(dsl -> 
Flux.from(dsl.select(BLOB_ID)
+        return postgresExecutor.executeRowsPaginated((dsl, lastRecord) -> 
dsl.select(BLOB_ID)
                 .from(TABLE_NAME)
                 .where(BUCKET_NAME.eq(bucketName.asString()))
-                .and(blobIdFrom.map(blobId -> 
BLOB_ID.greaterThan(blobId.asString())).orElseGet(DSL::noCondition))
-                .orderBy(BLOB_ID.asc())
-                .limit(batchSize)))
-            .map(record -> blobIdFactory.parse(record.get(BLOB_ID)))
-            .collectList()
-            .switchIfEmpty(Mono.just(ImmutableList.of()));
+                .and(lastRecord.map(record -> 
BLOB_ID.greaterThan(record.get(BLOB_ID))).orElseGet(DSL::noCondition))
+                .orderBy(BLOB_ID.asc()))
+            .map(record -> blobIdFactory.parse(record.get(BLOB_ID)));
     }
 
     private Hstore asHstore(BlobMetadata metadata) {
diff --git 
a/server/data/data-jmap-postgres/src/main/java/org/apache/james/jmap/postgres/upload/PostgresUploadDAO.java
 
b/server/data/data-jmap-postgres/src/main/java/org/apache/james/jmap/postgres/upload/PostgresUploadDAO.java
index 07f71e5188..35d4e22674 100644
--- 
a/server/data/data-jmap-postgres/src/main/java/org/apache/james/jmap/postgres/upload/PostgresUploadDAO.java
+++ 
b/server/data/data-jmap-postgres/src/main/java/org/apache/james/jmap/postgres/upload/PostgresUploadDAO.java
@@ -39,6 +39,7 @@ import org.apache.james.jmap.api.model.UploadId;
 import org.apache.james.jmap.api.model.UploadMetaData;
 import org.apache.james.mailbox.model.ContentType;
 import org.jooq.Record;
+import org.jooq.impl.DSL;
 
 import reactor.core.publisher.Flux;
 import reactor.core.publisher.Mono;
@@ -110,8 +111,10 @@ public class PostgresUploadDAO {
     }
 
     public Flux<Pair<UploadMetaData, Username>> 
listByUploadDateBefore(LocalDateTime before) {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.selectFrom(PostgresUploadTable.TABLE_NAME)
-                .where(PostgresUploadTable.UPLOAD_DATE.lessThan(before))))
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.selectFrom(PostgresUploadTable.TABLE_NAME)
+                .where(PostgresUploadTable.UPLOAD_DATE.lessThan(before))
+                .and(lastRecord.map(record -> 
PostgresUploadTable.ID.greaterThan(record.get(PostgresUploadTable.ID))).orElseGet(DSL::noCondition))
+                .orderBy(PostgresUploadTable.ID))
             .map(record -> Pair.of(uploadMetaDataFromRow(record), 
Username.of(record.get(PostgresUploadTable.USER_NAME))));
     }
 
diff --git 
a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java
 
b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java
index 4524e5dc67..462489df2b 100644
--- 
a/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java
+++ 
b/server/data/data-postgres/src/main/java/org/apache/james/mailrepository/postgres/PostgresMailRepositoryContentDAO.java
@@ -73,6 +73,7 @@ import org.apache.mailet.AttributeValue;
 import org.apache.mailet.Mail;
 import org.apache.mailet.PerRecipientHeaders;
 import org.jooq.Record;
+import org.jooq.impl.DSL;
 import org.jooq.postgres.extensions.types.Hstore;
 
 import com.fasterxml.jackson.databind.JsonNode;
@@ -210,9 +211,11 @@ public class PostgresMailRepositoryContentDAO {
     }
 
     private Flux<MailKey> listMailKeys(MailRepositoryUrl url) {
-        return postgresExecutor.executeRows(context -> 
Flux.from(context.select(KEY)
+        return postgresExecutor.executeRowsPaginated((context, lastRecord) -> 
context.select(KEY)
                 .from(TABLE_NAME)
-                .where(URL.eq(url.asString()))))
+                .where(URL.eq(url.asString()))
+                .and(lastRecord.map(record -> 
KEY.greaterThan(record.get(KEY))).orElseGet(DSL::noCondition))
+                .orderBy(KEY))
             .map(record -> new MailKey(record.get(KEY)));
     }
 
@@ -348,8 +351,10 @@ public class PostgresMailRepositoryContentDAO {
     }
 
     public Flux<BlobId> listBlobs() {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(HEADER_BLOB_ID, BODY_BLOB_ID)
-                .from(TABLE_NAME)))
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select(URL, KEY, HEADER_BLOB_ID, BODY_BLOB_ID)
+                .from(TABLE_NAME)
+                .where(lastRecord.map(record -> DSL.row(URL, 
KEY).greaterThan(record.get(URL), record.get(KEY))).orElseGet(DSL::noCondition))
+                .orderBy(URL, KEY))
             .flatMapIterable(record -> 
ImmutableList.of(blobIdFactory.parse(record.get(HEADER_BLOB_ID)), 
blobIdFactory.parse(record.get(BODY_BLOB_ID))));
     }
 }
diff --git 
a/server/data/data-postgres/src/main/java/org/apache/james/rrt/postgres/PostgresRecipientRewriteTableDAO.java
 
b/server/data/data-postgres/src/main/java/org/apache/james/rrt/postgres/PostgresRecipientRewriteTableDAO.java
index 07e050e698..8e767cfb7b 100644
--- 
a/server/data/data-postgres/src/main/java/org/apache/james/rrt/postgres/PostgresRecipientRewriteTableDAO.java
+++ 
b/server/data/data-postgres/src/main/java/org/apache/james/rrt/postgres/PostgresRecipientRewriteTableDAO.java
@@ -33,6 +33,7 @@ import org.apache.james.rrt.lib.Mapping;
 import org.apache.james.rrt.lib.MappingSource;
 import org.apache.james.rrt.lib.Mappings;
 import org.apache.james.rrt.lib.MappingsImpl;
+import org.jooq.impl.DSL;
 
 import com.google.common.collect.ImmutableList;
 
@@ -75,7 +76,11 @@ public class PostgresRecipientRewriteTableDAO {
     }
 
     public Flux<Pair<MappingSource, Mapping>> getAllMappings() {
-        return postgresExecutor.executeRows(dsl -> 
Flux.from(dsl.selectFrom(TABLE_NAME)))
+        return postgresExecutor.executeRowsPaginated((dsl, lastRecord) -> 
dsl.selectFrom(TABLE_NAME)
+                .where(lastRecord.map(record -> DSL.row(USERNAME, DOMAIN_NAME, 
TARGET_ADDRESS)
+                        .greaterThan(record.get(USERNAME), 
record.get(DOMAIN_NAME), record.get(TARGET_ADDRESS)))
+                    .orElseGet(DSL::noCondition))
+                .orderBy(USERNAME, DOMAIN_NAME, TARGET_ADDRESS))
             .map(record -> Pair.of(
                 MappingSource.fromUser(record.get(USERNAME), 
record.get(DOMAIN_NAME)),
                 Mapping.of(record.get(TARGET_ADDRESS))));
diff --git 
a/server/data/data-postgres/src/main/java/org/apache/james/user/postgres/PostgresUsersDAO.java
 
b/server/data/data-postgres/src/main/java/org/apache/james/user/postgres/PostgresUsersDAO.java
index 908f086462..01e0b60542 100644
--- 
a/server/data/data-postgres/src/main/java/org/apache/james/user/postgres/PostgresUsersDAO.java
+++ 
b/server/data/data-postgres/src/main/java/org/apache/james/user/postgres/PostgresUsersDAO.java
@@ -20,7 +20,6 @@
 package org.apache.james.user.postgres;
 
 import static 
org.apache.james.backends.postgres.utils.PostgresExecutor.DEFAULT_INJECT;
-import static 
org.apache.james.backends.postgres.utils.PostgresExecutor.EAGER_FETCH;
 import static 
org.apache.james.user.postgres.PostgresUserDataDefinition.PostgresUserTable.ALGORITHM;
 import static 
org.apache.james.user.postgres.PostgresUserDataDefinition.PostgresUserTable.AUTHORIZED_USERS;
 import static 
org.apache.james.user.postgres.PostgresUserDataDefinition.PostgresUserTable.DELEGATED_USERS;
@@ -46,6 +45,7 @@ import org.apache.james.user.api.model.User;
 import org.apache.james.user.lib.UsersDAO;
 import org.apache.james.user.lib.model.Algorithm;
 import org.apache.james.user.lib.model.DefaultUser;
+import org.jooq.Condition;
 import org.jooq.DSLContext;
 import org.jooq.Field;
 import org.jooq.Record;
@@ -142,9 +142,7 @@ public class PostgresUsersDAO implements UsersDAO {
 
     @Override
     public Flux<Username> listReactive() {
-        return postgresExecutor.executeRows(dslContext -> 
Flux.from(dslContext.select(USERNAME)
-                .from(TABLE_NAME)), EAGER_FETCH)
-            .map(record -> Username.of(record.get(USERNAME)));
+        return listUsernames(DSL.noCondition());
     }
 
     @Override
@@ -153,10 +151,15 @@ public class PostgresUsersDAO implements UsersDAO {
             return listReactive();
         }
         String domainPattern = "%@" + domain.asString();
-        return postgresExecutor.executeRows(dslContext -> Flux.from(
-                dslContext.select(USERNAME)
-                    .from(TABLE_NAME)
-                    .where(USERNAME.like(domainPattern))), EAGER_FETCH)
+        return listUsernames(USERNAME.like(domainPattern));
+    }
+
+    private Flux<Username> listUsernames(Condition condition) {
+        return postgresExecutor.executeRowsPaginated((dslContext, lastRecord) 
-> dslContext.select(USERNAME)
+                .from(TABLE_NAME)
+                .where(condition)
+                .and(lastRecord.map(record -> 
USERNAME.greaterThan(record.get(USERNAME))).orElseGet(DSL::noCondition))
+                .orderBy(USERNAME))
             .map(record -> Username.of(record.get(USERNAME)));
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to