Skip to content

Commit ba7c996

Browse files
committed
Use PIT + search_after for delete in serverless mode (#691)
Signed-off-by: Sotaro Hikita <bering1814@gmail.com>
1 parent 9807b25 commit ba7c996

1 file changed

Lines changed: 20 additions & 9 deletions

File tree

mr/src/main/java/org/opensearch/hadoop/rest/RestRepository.java

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -438,14 +438,6 @@ public void delete() {
438438
// 250 results
439439

440440
int batchSize = 500;
441-
StringBuilder sb = new StringBuilder(resources.getResourceWrite().index());
442-
if (resources.getResourceWrite().isTyped()) {
443-
sb.append('/').append(resources.getResourceWrite().type());
444-
}
445-
sb.append("/_search?scroll=10m&_source=false&size=");
446-
sb.append(batchSize);
447-
sb.append("&sort=_doc");
448-
String scanQuery = sb.toString();
449441
ScrollReader scrollReader = new ScrollReader(
450442
ScrollReaderConfigBuilder.builder(new JdkValueReader(), settings)
451443
.setReadMetadata(true)
@@ -458,8 +450,27 @@ public void delete() {
458450
.setErrorHandlerLoader(new AbortOnlyHandlerLoader()) // Only abort since this is internal
459451
);
460452

453+
ScrollQuery sq;
454+
if (this.settings.getServerlessMode()) {
455+
// Serverless mode: use PIT + search_after instead of scroll API
456+
String searchQuery = "_search?size=" + batchSize + "&_source=false&track_total_hits=true";
457+
BytesArray body = new BytesArray("{\"query\":{\"match_all\":{}},\"sort\":[\"_doc\",\"_id\"]}");
458+
String index = resources.getResourceWrite().index();
459+
String keepAlive = settings.getPitKeepAlive();
460+
sq = scanLimitSearchAfter(searchQuery, body, -1, scrollReader, index, keepAlive);
461+
} else {
462+
StringBuilder sb = new StringBuilder(resources.getResourceWrite().index());
463+
if (resources.getResourceWrite().isTyped()) {
464+
sb.append('/').append(resources.getResourceWrite().type());
465+
}
466+
sb.append("/_search?scroll=10m&_source=false&size=");
467+
sb.append(batchSize);
468+
sb.append("&sort=_doc");
469+
String scanQuery = sb.toString();
470+
sq = scanAll(scanQuery, null, scrollReader);
471+
}
472+
461473
// start iterating
462-
ScrollQuery sq = scanAll(scanQuery, null, scrollReader);
463474
try {
464475
BytesArray entry = new BytesArray(0);
465476

0 commit comments

Comments
 (0)