Skip to content

Commit a82ba89

Browse files
committed
fix(hstore): scope ordered range scan to indexes
1 parent 7c85cca commit a82ba89

14 files changed

Lines changed: 855 additions & 38 deletions

File tree

hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/page/IdHolder.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -131,7 +131,7 @@ public PageIds fetchNext(String page, long pageSize) {
131131

132132
PageIds result = this.fetcher.apply((ConditionQuery) this.query);
133133
assert result != null;
134-
if (result.ids().size() < pageSize || result.page() == null) {
134+
if (result.page() == null) {
135135
this.exhausted = true;
136136
}
137137
return result;

hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/page/PageEntryIterator.java

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -87,7 +87,12 @@ private boolean fetch() {
8787
this.remaining -= this.pageResults.total();
8888
return true;
8989
} else {
90-
this.pageInfo.increase();
90+
if (this.pageResults.hasNextPage() &&
91+
!this.pageResults.page().equals(this.pageInfo.page())) {
92+
this.pageInfo.page(this.pageResults.page());
93+
} else {
94+
this.pageInfo.increase();
95+
}
9196
return this.fetch();
9297
}
9398
}

hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/page/QueryList.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -272,7 +272,8 @@ public PageResults<R> iterator(int index, String page, long pageSize) {
272272
this.updateResultsFilter(bindQuery);
273273
PageIds pageIds = holder.fetchNext(page, pageSize);
274274
if (pageIds.empty()) {
275-
return PageResults.emptyIterator();
275+
return PageResults.emptyIterator(bindQuery,
276+
pageIds.pageState());
276277
}
277278

278279
QueryResults<R> results = this.queryByIndexIds(pageIds.ids(),
@@ -365,5 +366,12 @@ public long total() {
365366
public static <R> PageResults<R> emptyIterator() {
366367
return (PageResults<R>) EMPTY;
367368
}
369+
370+
public static <R> PageResults<R> emptyIterator(Query query,
371+
PageState pageState) {
372+
return new PageResults<>(
373+
new QueryResults<>(QueryResults.emptyIterator(), query),
374+
pageState);
375+
}
368376
}
369377
}

hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/tx/GraphIndexTransaction.java

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -749,8 +749,15 @@ private PageIds doIndexQueryOnce(IndexLabel indexLabel,
749749
Query.checkForceCapacity(ids.size());
750750
this.recordIndexValue(query, index);
751751
}
752-
// If there is no data, the entries is not a Metadatable object
753752
if (ids.isEmpty()) {
753+
if (query.paging() && entries instanceof Metadatable) {
754+
PageState pageState = PageInfo.pageState(entries);
755+
if (pageState.position().length > 0) {
756+
pageState = new PageState(pageState.position(),
757+
pageState.offset(), 0);
758+
return new PageIds(ids, pageState);
759+
}
760+
}
754761
return PageIds.EMPTY;
755762
}
756763
// NOTE: Memory backend's iterator is not Metadatable
@@ -761,7 +768,10 @@ private PageIds doIndexQueryOnce(IndexLabel indexLabel,
761768
"The entries must be Metadatable when query " +
762769
"in paging, but got '%s'",
763770
entries.getClass().getName());
764-
return new PageIds(ids, PageInfo.pageState(entries));
771+
PageState pageState = PageInfo.pageState(entries);
772+
pageState = new PageState(pageState.position(),
773+
pageState.offset(), ids.size());
774+
return new PageIds(ids, pageState);
765775
} finally {
766776
locks.unlock();
767777
CloseableIterator.closeIterator(entries);

hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessions.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,15 @@ public BackendColumnIterator scan(String table,
162162
scanType, query);
163163
}
164164

165+
public abstract BackendColumnIterator scanOrdered(String table,
166+
byte[] ownerKeyFrom,
167+
byte[] ownerKeyTo,
168+
byte[] keyFrom,
169+
byte[] keyTo,
170+
int scanType,
171+
byte[] query,
172+
long limit);
173+
165174
public abstract BackendColumnIterator scan(String table,
166175
byte[] ownerKeyFrom,
167176
byte[] ownerKeyTo,

hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessionsImpl.java

Lines changed: 34 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@
5050
import org.apache.hugegraph.store.HgScanQuery;
5151
import org.apache.hugegraph.store.HgStoreClient;
5252
import org.apache.hugegraph.store.HgStoreSession;
53+
import org.apache.hugegraph.store.client.NodeTxSessionProxy;
5354
import org.apache.hugegraph.store.client.grpc.KvCloseableIterator;
5455
import org.apache.hugegraph.store.client.util.HgStoreClientConst;
5556
import org.apache.hugegraph.store.grpc.common.ScanOrderType;
@@ -226,22 +227,30 @@ private static class ColumnIterator<T extends HgKvIterator> implements
226227
private final int scanType;
227228
private final String table;
228229
private final byte[] value;
230+
private final boolean keepPositionAfterExhausted;
229231
private boolean gotNext;
230232
private byte[] position;
231233

232234
public ColumnIterator(String table, T results) {
233-
this(table, results, null, null, 0);
235+
this(table, results, null, null, 0, false);
234236
}
235237

236238
public ColumnIterator(String table, T results, byte[] keyBegin,
237239
byte[] keyEnd, int scanType) {
240+
this(table, results, keyBegin, keyEnd, scanType, false);
241+
}
242+
243+
public ColumnIterator(String table, T results, byte[] keyBegin,
244+
byte[] keyEnd, int scanType,
245+
boolean keepPositionAfterExhausted) {
238246
E.checkNotNull(results, "results");
239247
this.table = table;
240248
this.iter = results;
241249
this.keyBegin = keyBegin;
242250
this.keyEnd = keyEnd;
243251
this.scanType = scanType;
244252
this.value = null;
253+
this.keepPositionAfterExhausted = keepPositionAfterExhausted;
245254
if (this.iter.hasNext()) {
246255
this.iter.next();
247256
this.gotNext = true;
@@ -317,11 +326,9 @@ private boolean match(int expected) {
317326

318327
@Override
319328
public boolean hasNext() {
320-
if (gotNext) {
329+
if (gotNext && !this.keepPositionAfterExhausted) {
321330
this.position = this.iter.position();
322-
} else {
323-
// QUESTION: Resetting the position may result in the caller being unable to
324-
// retrieve the corresponding position.
331+
} else if (!this.keepPositionAfterExhausted) {
325332
this.position = null;
326333
}
327334
return gotNext;
@@ -376,6 +383,9 @@ public BackendColumn next() {
376383
BackendColumn col =
377384
BackendColumn.of(this.iter.key(),
378385
this.iter.value());
386+
if (this.keepPositionAfterExhausted) {
387+
this.position = col.name;
388+
}
379389
if (this.iter.hasNext()) {
380390
gotNext = true;
381391
this.iter.next();
@@ -716,6 +726,25 @@ public BackendColumnIterator scan(String table, byte[] ownerKeyFrom,
716726
scanType);
717727
}
718728

729+
@Override
730+
public BackendColumnIterator scanOrdered(String table,
731+
byte[] ownerKeyFrom,
732+
byte[] ownerKeyTo,
733+
byte[] keyFrom,
734+
byte[] keyTo,
735+
int scanType,
736+
byte[] query,
737+
long limit) {
738+
assert !this.hasChanges();
739+
HgKvIterator<HgKvEntry> result =
740+
((NodeTxSessionProxy) this.graph).scanIteratorOrdered(
741+
table, HgOwnerKey.of(ownerKeyFrom, keyFrom),
742+
HgOwnerKey.of(ownerKeyTo, keyTo),
743+
toHstoreLimit(limit), scanType, query);
744+
return new ColumnIterator<>(table, result, keyFrom, keyTo,
745+
scanType, true);
746+
}
747+
719748
@Override
720749
public BackendColumnIterator scan(String table, byte[] ownerKeyFrom,
721750
byte[] ownerKeyTo,

hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreTable.java

Lines changed: 64 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -114,15 +114,12 @@ private static boolean direction(Condition condition) {
114114

115115
protected static BackendEntryIterator newEntryIterator(
116116
BackendColumnIterator cols, Query query) {
117-
return new BinaryEntryIterator<>(cols, query, (entry, col) -> {
118-
if (entry == null || !entry.belongToMe(col)) {
119-
HugeType type = query.resultType();
120-
// NOTE: only support BinaryBackendEntry currently
121-
entry = new BinaryBackendEntry(type, col.name);
122-
}
123-
entry.columns(col);
124-
return entry;
125-
});
117+
BiFunction<BackendEntry, BackendColumn, BackendEntry> merger =
118+
(entry, col) -> mergeColumn(query, entry, col);
119+
if (query.resultType().isRangeIndex()) {
120+
return new RangeIndexEntryIterator(cols, query, merger);
121+
}
122+
return new BinaryEntryIterator<>(cols, query, merger);
126123
}
127124

128125
protected static BackendEntryIterator newEntryIteratorOlap(
@@ -138,6 +135,42 @@ protected static BackendEntryIterator newEntryIteratorOlap(
138135
});
139136
}
140137

138+
private static BackendEntry mergeColumn(Query query, BackendEntry entry,
139+
BackendColumn col) {
140+
if (entry == null || !entry.belongToMe(col)) {
141+
HugeType type = query.resultType();
142+
// NOTE: only support BinaryBackendEntry currently
143+
entry = new BinaryBackendEntry(type, col.name);
144+
}
145+
entry.columns(col);
146+
return entry;
147+
}
148+
149+
private static final class RangeIndexEntryIterator
150+
extends BinaryEntryIterator<BackendColumn> {
151+
152+
private byte[] lastPosition;
153+
154+
private RangeIndexEntryIterator(
155+
BackendColumnIterator cols, Query query,
156+
BiFunction<BackendEntry, BackendColumn, BackendEntry> merger) {
157+
super(cols, query, merger);
158+
this.lastPosition = PageState.EMPTY_BYTES;
159+
}
160+
161+
@Override
162+
public BackendEntry next() {
163+
BackendEntry entry = super.next();
164+
this.lastPosition = entry.id().asBytes();
165+
return entry;
166+
}
167+
168+
@Override
169+
protected PageState pageState() {
170+
return new PageState(this.lastPosition, 0, (int) this.count());
171+
}
172+
}
173+
141174
public static String bytes2String(byte[] bytes) {
142175
StringBuilder result = new StringBuilder();
143176
for (byte b : bytes) {
@@ -625,17 +658,26 @@ protected BackendColumnIterator queryByRange(Session session,
625658
}
626659
ConditionQuery cq;
627660
Query origin = query.originQuery();
628-
if (query.paging() && !query.page().isEmpty()) {
629-
type = (type & ~Session.SCAN_GTE_BEGIN) | Session.SCAN_GT_BEGIN;
630-
}
661+
byte[] position = null;
631662
byte[] ownerStart = this.ownerByQueryDelegate.apply(query.resultType(),
632663
query.start());
633664
byte[] ownerEnd = this.ownerByQueryDelegate.apply(query.resultType(),
634665
query.end());
666+
if (query.resultType().isRangeIndex()) {
667+
if (query.paging() && !query.page().isEmpty()) {
668+
start = rangeIndexScanStart(query, start);
669+
type = (type & ~Session.SCAN_GTE_BEGIN) | Session.SCAN_GT_BEGIN;
670+
}
671+
return session.scanOrdered(this.table(), ownerStart, ownerEnd,
672+
start, end, type, null,
673+
rangeScanLimit(query));
674+
}
675+
if (query.paging() && !query.page().isEmpty()) {
676+
position = PageState.fromString(query.page()).position();
677+
}
635678
if (origin instanceof ConditionQuery &&
636679
(query.resultType().isEdge() || query.resultType().isVertex())) {
637680
cq = (ConditionQuery) query.originQuery();
638-
long limit = rangeScanLimit(query);
639681

640682
// LOG.debug("query {} with ownerKeyFrom: {}, ownerKeyTo: {}, " +
641683
// "keyFrom: {}, keyTo: {}, " +
@@ -644,10 +686,17 @@ protected BackendColumnIterator queryByRange(Session session,
644686
// bytes2String(ownerEnd), bytes2String(start),
645687
// bytes2String(end), type, cq.bytes());
646688
return session.scan(this.table(), ownerStart, ownerEnd, start,
647-
end, type, cq.bytes(), null, limit);
689+
end, type, cq.bytes(), position);
648690
}
649691
return session.scan(this.table(), ownerStart, ownerEnd, start, end,
650-
type, null, null, rangeScanLimit(query));
692+
type, null, position);
693+
}
694+
695+
static byte[] rangeIndexScanStart(IdRangeQuery query, byte[] start) {
696+
if (query.paging() && !query.page().isEmpty()) {
697+
return PageState.fromString(query.page()).position();
698+
}
699+
return start;
651700
}
652701

653702
private static long rangeScanLimit(IdRangeQuery query) {

0 commit comments

Comments
 (0)