From b9b857b0b1f3f9df35ea7ecf9e88e1e5550bd43c Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Wed, 12 Aug 2026 18:20:57 +0000
Subject: [PATCH] Serve Cassandra cursor repositioning with CQL slice queries instead of a no-op iterator restart (#865)
---
opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java | 243 ++++++++++++++++++++++++++++++++++++++----------
1 files changed, 190 insertions(+), 53 deletions(-)
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java b/opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java
index 7efeb8c..e30c748 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java
@@ -58,6 +58,7 @@
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.config.DriverConfigLoader;
+import com.datastax.oss.driver.api.core.cql.BoundStatement;
import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import com.datastax.oss.driver.api.core.cql.ResultSet;
import com.datastax.oss.driver.api.core.cql.Row;
@@ -71,7 +72,18 @@
private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
private CASBackendCfg config;
-
+
+ //a cursor starts with a small page so that a point lookup does not transfer thousands of rows,
+ //and doubles it while a scan keeps outrunning it (the driver default page is 5000 rows)
+ final int initialPageSize=Math.max(1,Integer.getInteger("org.openidentityplatform.opendj.cassandra.fetchsize.initial",32));
+ final int maxPageSize=Math.max(initialPageSize,Integer.getInteger("org.openidentityplatform.opendj.cassandra.fetchsize",1000));
+
+ //CursorImpl.forwardInPage outcomes
+ final static int POSITIONED=1;
+ final static int MISSING=0;
+ final static int NEEDS_SEEK=-1;
+ final static int PAGE_OUT=-2;
+
public CASStorage(CASBackendCfg cfg, ServerContext serverContext) {
this.config = cfg;
cfg.addCASChangeListener(this);
@@ -292,42 +304,138 @@
}
}
- private final class CursorImpl implements Cursor<ByteString, ByteString> {
+ // Iterates the (baseDN,indexId) partition in key order.
+ // A ResultSet can only be consumed once: rc.iterator() always returns the same iterator,
+ // so "restarting" it never rewinds. Every repositioning that cannot be served by moving
+ // forward within the already-fetched rows therefore runs a new server-side slice query
+ // on the "key" clustering column instead.
+ // Pages are sized by the cursor rather than by the driver default of 5000 rows: a point
+ // lookup must not transfer thousands of entries, so a cursor starts with a small page and
+ // doubles it (up to maxPageSize) every time a scan outruns it. Scanning past the end of a
+ // page runs the next slice from the current key instead of letting the driver page with the
+ // size the cursor started with.
+ final class CursorImpl implements Cursor<ByteString, ByteString> {
final TreeName treeName;
final TransactionImpl tx;
+ final String tableName=getTableName();
+
+ //visible for tests
+ long queryCount;
+ int pageSize=initialPageSize;
ResultSet rc;
Iterator<Row> iterator;
+ //the row the cursor is on, or - when defined is false - the row a failed positionToKey
+ //stopped just before, which is where next() resumes (same as pdb)
Row current=null;
-
+ boolean defined=false;
+ //rc holds a DESC page: it must never serve a forward move
+ boolean descending=false;
+ boolean closed=false;
+
public CursorImpl(TransactionImpl tx,TreeName treeName) {
this.treeName=treeName;
this.tx=tx;
- rc=full();
- iterator=rc.iterator();
+ //lazy: the first navigation decides which query to run, a seek must not pay for a partition scan
}
- ResultSet full(){
- return execute(
- prepared.get("SELECT key,value FROM "+getTableName()+" WHERE baseDN=:baseDN and indexId=:indexId ORDER BY key").bind()
- .setString("baseDN", treeName.getBaseDN()).setString("indexId", treeName.getIndexId())
- );
+ ResultSet select(String condition,ByteSequence key,int rows){
+ queryCount++;
+ BoundStatement statement=prepared.get("SELECT key,value FROM "+tableName+" WHERE baseDN=:baseDN and indexId=:indexId"+condition).bind()
+ .setString("baseDN", treeName.getBaseDN()).setString("indexId", treeName.getIndexId())
+ .setPageSize(rows);
+ if (key!=null) {
+ statement=statement.setByteBuffer("key", ByteBuffer.wrap(key.toByteArray()));
+ }
+ return execute(statement);
}
-
+
+ //runs a new page and positions the cursor on its first row
+ boolean slice(String condition,ByteSequence key){
+ rc=select(condition,key,pageSize);
+ iterator=rc.iterator();
+ descending=false;
+ if (iterator.hasNext()) {
+ current=iterator.next();
+ defined=true;
+ return true;
+ }
+ current=null;
+ defined=false;
+ return false;
+ }
+
+ boolean seek(ByteSequence key){
+ return slice(" and key>=:key ORDER BY key",key);
+ }
+
+ void growPage() {
+ pageSize=(int)Math.min(maxPageSize,2L*pageSize);
+ }
+
+ //serves a forward repositioning from the rows the driver already fetched: they are the
+ //sorted rows following the current one (blob clustering collates in unsigned byte order)
+ int forwardInPage(ByteSequence key,boolean exactMatch) {
+ if (!defined || descending || rc==null) {
+ return NEEDS_SEEK;
+ }
+ int cmp=key.compareTo(getKey());
+ if (cmp==0) {
+ return POSITIONED;
+ }
+ if (cmp<0) { //backward: only the server can rewind
+ return NEEDS_SEEK;
+ }
+ while (rc.getAvailableWithoutFetching()>0) {
+ current=iterator.next();
+ cmp=key.compareTo(getKey());
+ if (cmp==0) {
+ return POSITIONED;
+ }
+ if (cmp<0) { //walked past the key
+ if (exactMatch) {
+ defined=false; //the cursor stops just before this row
+ return MISSING;
+ }
+ return POSITIONED;
+ }
+ }
+ return PAGE_OUT;
+ }
+
@Override
public boolean next() {
- try {
- current=iterator.next();
- return true;
- }catch (NoSuchElementException e) {
- current=null;
+ if (closed) {
+ return false;
}
- return false;
+ if (current!=null && !defined) { //a failed positionToKey stopped just before this row
+ defined=true;
+ return true;
+ }
+ if (rc==null) { //lazy cursor: the first navigation runs the scan
+ return slice(" ORDER BY key",null);
+ }
+ if (!descending && rc.getAvailableWithoutFetching()>0) {
+ current=iterator.next();
+ defined=true;
+ return true;
+ }
+ if (current==null) { //exhausted or explicitly undefined: stay there
+ return false;
+ }
+ if (!descending && rc.isFullyFetched()) { //the server has nothing left either
+ current=null;
+ defined=false;
+ return false;
+ }
+ //page boundary: continue with a bigger slice starting right after the current key
+ growPage();
+ return slice(" and key>:key ORDER BY key",getKey());
}
@Override
public boolean isDefined() {
- return current!=null;
+ return defined;
}
@Override
@@ -354,72 +462,101 @@
tx.delete(treeName, getKey());
}
+ //a closed cursor is undefined and every navigation on it returns false, like EmptyCursor
@Override
public void close() {
- iterator=null;
+ closed=true;
+ iterator=Collections.emptyIterator();
current=null;
+ defined=false;
+ descending=false;
rc=null;
}
@Override
public boolean positionToKeyOrNext(ByteSequence key) {
- if (!isDefined() || key.compareTo(getKey())<0) { //restart iterator
- iterator=rc.iterator();
+ if (closed) {
+ return false;
}
- while (iterator.hasNext()) {
- current=iterator.next();
- if (key.compareTo(getKey())<=0) {
- return true;
- }
+ final int served=forwardInPage(key,false);
+ if (served>=0) {
+ return served==POSITIONED;
}
- current=null;
- return false;
+ if (served==PAGE_OUT) { //the scan outran its page: the next slice should be bigger
+ growPage();
+ }
+ return seek(key);
}
-
+
@Override
public boolean positionToKey(ByteSequence key) {
- if (!isDefined() || key.compareTo(getKey())<0) { //restart iterator
- iterator=rc.iterator();
+ if (closed) {
+ return false;
}
- if (isDefined() && key.compareTo(getKey())==0) {
+ final int served=forwardInPage(key,true);
+ if (served>=0) {
+ return served==POSITIONED;
+ }
+ if (served==PAGE_OUT) {
+ growPage();
+ }
+ if (seek(key) && key.compareTo(getKey())==0) {
return true;
}
- while (iterator.hasNext()) {
- current=iterator.next();
- if (key.compareTo(getKey())==0) {
- return true;
- }
- }
- current=null;
+ //like jeb/pdb a miss leaves the cursor undefined; the row the seek landed on is the
+ //first key after the missing one, so next() resumes there instead of skipping it
+ defined=false;
return false;
}
-
+
@Override
public boolean positionToLastKey() {
- while (iterator.hasNext()) {
- current=iterator.next();
+ if (closed) {
+ return false;
}
- if (current!=null) {
+ rc=select(" ORDER BY key DESC LIMIT 1",null,1);
+ iterator=rc.iterator();
+ descending=true;
+ if (iterator.hasNext()) {
+ current=iterator.next(); //nothing follows the last key
+ defined=true;
return true;
}
+ current=null;
+ defined=false;
return false;
}
@Override
public boolean positionToIndex(int index) {
- iterator=rc.iterator(); //restart iterator
- int ct=0;
- while(iterator.hasNext()){
- current=iterator.next();
- if (ct==index) {
- return true;
- }
- ct++;
+ if (closed) {
+ return false;
}
- current=null;
- return false;
+ if (index<0) {
+ rc=null; //an invalid index resets the cursor: next() starts the scan again
+ iterator=null;
+ current=null;
+ defined=false;
+ descending=false;
+ return false;
+ }
+ //CQL has no offset clause: restart from the first row and skip. The rows that have to
+ //be walked are asked for in one page instead of one round trip per page
+ rc=select(" ORDER BY key",null,(int)Math.min(maxPageSize,Math.max(pageSize,index+1L)));
+ iterator=rc.iterator();
+ descending=false;
+ for (int ct=0;ct<=index;ct++) {
+ if (!iterator.hasNext()) {
+ current=null;
+ defined=false;
+ return false;
+ }
+ current=iterator.next();
+ }
+ defined=true;
+ return true;
}
}
--
Gitblit v1.10.0