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]