Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ Inspired from [Keep a Changelog](https://keepachangelog.com/en/1.0.0/)
- Fixed opensearch-spark-40 dependency license SHAs and added dependabot config for sql-40 module ([#688](https://github.com/opensearch-project/opensearch-hadoop/pull/688))
- Fixed global refresh when using dynamic index patterns ([#686](https://github.com/opensearch-project/opensearch-hadoop/pull/686))
- Fixed build failures when downloading Apache project dependencies (Hadoop, Hive, Spark) ([#595](https://github.com/opensearch-project/opensearch-hadoop/pull/595))
- Fixed serverless mode SaveMode.Overwrite failing when document count exceeds scroll size ([#693](https://github.com/opensearch-project/opensearch-hadoop/pull/693))

### Security

Expand Down
29 changes: 20 additions & 9 deletions mr/src/main/java/org/opensearch/hadoop/rest/RestRepository.java
Original file line number Diff line number Diff line change
Expand Up @@ -438,14 +438,6 @@ public void delete() {
// 250 results

int batchSize = 500;
StringBuilder sb = new StringBuilder(resources.getResourceWrite().index());
if (resources.getResourceWrite().isTyped()) {
sb.append('/').append(resources.getResourceWrite().type());
}
sb.append("/_search?scroll=10m&_source=false&size=");
sb.append(batchSize);
sb.append("&sort=_doc");
String scanQuery = sb.toString();
ScrollReader scrollReader = new ScrollReader(
ScrollReaderConfigBuilder.builder(new JdkValueReader(), settings)
.setReadMetadata(true)
Expand All @@ -458,8 +450,27 @@ public void delete() {
.setErrorHandlerLoader(new AbortOnlyHandlerLoader()) // Only abort since this is internal
);

ScrollQuery sq;
if (this.settings.getServerlessMode()) {
// Serverless mode: use PIT + search_after instead of scroll API
String searchQuery = "_search?size=" + batchSize + "&_source=false&track_total_hits=true";
BytesArray body = new BytesArray("{\"query\":{\"match_all\":{}},\"sort\":[\"_doc\",\"_id\"]}");
String index = resources.getResourceWrite().index();
String keepAlive = settings.getPitKeepAlive();
sq = scanLimitSearchAfter(searchQuery, body, -1, scrollReader, index, keepAlive);
} else {
StringBuilder sb = new StringBuilder(resources.getResourceWrite().index());
if (resources.getResourceWrite().isTyped()) {
sb.append('/').append(resources.getResourceWrite().type());
}
sb.append("/_search?scroll=10m&_source=false&size=");
sb.append(batchSize);
sb.append("&sort=_doc");
String scanQuery = sb.toString();
sq = scanAll(scanQuery, null, scrollReader);
}

// start iterating
ScrollQuery sq = scanAll(scanQuery, null, scrollReader);
try {
BytesArray entry = new BytesArray(0);

Expand Down
Loading