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 .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -70,3 +70,4 @@ testfixtures_shared/
# build files generated
doc-tools/missing-doclet/bin/
/sandbox/plugins/engine-datafusion/target/
**/Cargo.lock
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
- Add support for enabling pluggable data formats, starting with phase-1 of decoupling shard from engine, and introducing basic abstractions ([#20675](https://github.com/opensearch-project/OpenSearch/pull/20675))
- Add concurrent queue in libs and composite engine sandbox plugin ([#20909](https://github.com/opensearch-project/OpenSearch/pull/20909))
- Add interface for the Multi format merge flow ([#20908](https://github.com/opensearch-project/OpenSearch/pull/20908))
- Add CatalogSnapshotManager lifecycle management with reference-counted snapshot tracking and serialization support for Segment and WriterFileSet ([#20982](https://github.com/opensearch-project/OpenSearch/pull/20982))

- Add warmup phase to wait for lag to catch up in pull-based ingestion before serving ([#20526](https://github.com/opensearch-project/OpenSearch/pull/20526))
- Add a new static method to IndicesOptions API to expose `STRICT_EXPAND_OPEN_HIDDEN_FORBID_CLOSED` index option ([#20980](https://github.com/opensearch-project/OpenSearch/pull/20980))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@

import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.EngineReaderManager;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.shard.ShardPath;

import java.io.IOException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@
import org.apache.lucene.search.ReferenceManager;
import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.EngineReaderManager;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.io.IOException;
import java.util.Collection;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,15 +30,17 @@
import org.opensearch.cluster.metadata.IndexMetadata;
import org.opensearch.cluster.metadata.Metadata;
import org.opensearch.cluster.service.ClusterService;
import org.opensearch.common.concurrent.GatedCloseable;
import org.opensearch.core.index.Index;
import org.opensearch.index.IndexService;
import org.opensearch.index.engine.DataFormatAwareEngine;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.dataformat.FieldTypeCapabilities;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.EngineReaderManager;
import org.opensearch.index.engine.exec.Segment;
import org.opensearch.index.engine.exec.WriterFileSet;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.engine.exec.coord.CatalogSnapshotManager;
import org.opensearch.index.shard.IndexShard;
import org.opensearch.indices.IndicesService;
import org.opensearch.test.OpenSearchTestCase;
Expand Down Expand Up @@ -104,13 +106,15 @@ public void testEndToEndExecuteWithMockBackend() throws IOException {

Segment seg1 = Segment.builder(0L).addSearchableFiles(format, fs1).build();
Segment seg2 = Segment.builder(1L).addSearchableFiles(format, fs2).build();
MockCatalogSnapshot snapshot = new MockCatalogSnapshot(1L, List.of(seg1, seg2), format);

CatalogSnapshotManager snapshotManager = new CatalogSnapshotManager(1L, 1L, 0L, List.of(seg1, seg2), 2L, Map.of());

MockReaderManager readerManager = new MockReaderManager(format.name());
readerManager.afterRefresh(true, snapshot);
try (GatedCloseable<CatalogSnapshot> ref = snapshotManager.acquireSnapshot()) {
readerManager.afterRefresh(true, ref.get());
}

DataFormatAwareEngine engine = new DataFormatAwareEngine(Map.of(format, readerManager));
engine.setLatestSnapshot(snapshot);
DataFormatAwareEngine engine = new DataFormatAwareEngine(Map.of(format, readerManager), snapshotManager);

// Mock shard + cluster wiring
IndexShard shard = mock(IndexShard.class);
Expand Down Expand Up @@ -288,16 +292,18 @@ public String serializeToString() {
}

@Override
public void setCatalogSnapshotMap(Map<Long, ? extends CatalogSnapshot> map) {}

@Override
public void setUserData(Map<String, String> userData, boolean b) {}
public void setUserData(Map<String, String> userData) {}

@Override
public Object getReader(DataFormat dataFormat) {
return null;
}

@Override
public MockCatalogSnapshot clone() {
return new MockCatalogSnapshot(generation, segments, format);
}

@Override
protected void closeInternal() {}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,11 @@
import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.common.concurrent.GatedCloseable;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.DataFormatAwareEngineFactory;
import org.opensearch.index.engine.exec.EngineReaderManager;
import org.opensearch.index.engine.exec.IndexReaderProvider;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.engine.exec.coord.CatalogSnapshotManager;

import java.io.Closeable;
import java.io.IOException;
Expand All @@ -35,62 +36,56 @@
public class DataFormatAwareEngine implements IndexReaderProvider, Closeable {

private final Map<DataFormat, EngineReaderManager<?>> readerManagers;
private volatile CatalogSnapshot latestSnapshot;
private volatile CatalogSnapshotManager catalogSnapshotManager;

/**
* Constructs a new DataFormatAwareEngine with pre-built maps.
* Constructs a new DataFormatAwareEngine.
* Prefer using {@link DataFormatAwareEngineFactory#create()}.
*/
public DataFormatAwareEngine(Map<DataFormat, EngineReaderManager<?>> readerManagers, CatalogSnapshotManager catalogSnapshotManager) {
this.readerManagers = readerManagers;
this.catalogSnapshotManager = catalogSnapshotManager;
}

/**
* Constructs a new DataFormatAwareEngine without a snapshot manager.
* The manager must be set via {@link #setCatalogSnapshotManager} before acquiring readers.
*/
public DataFormatAwareEngine(Map<DataFormat, EngineReaderManager<?>> readerManagers) {
this.readerManagers = readerManagers;
}

public EngineReaderManager<?> getReaderManager(DataFormat format) {
return readerManagers.get(format);
public void setCatalogSnapshotManager(CatalogSnapshotManager catalogSnapshotManager) {
this.catalogSnapshotManager = catalogSnapshotManager;
}

/**
* Called by the catalog snapshot lifecycle listener after a refresh
* to update the latest searchable snapshot.
*/
public void setLatestSnapshot(CatalogSnapshot snapshot) {
CatalogSnapshot prev = this.latestSnapshot;
this.latestSnapshot = snapshot;
if (prev != null) {
prev.decRef();
}
public EngineReaderManager<?> getReaderManager(DataFormat format) {
return readerManagers.get(format);
}

/**
* Acquires a DataFormatAwareReader on the latest catalog snapshot.
* The snapshot is incRef'd; the caller MUST close the returned
* {@link DataFormatAwareReader} when done, which decRef's the snapshot.
* The caller MUST close the returned {@link DataFormatAwareReader} when done,
* which releases the snapshot reference.
*/
public GatedCloseable<Reader> acquireReader() throws IOException {
CatalogSnapshot snapshot = latestSnapshot;
if (snapshot == null) {
throw new IllegalStateException("No catalog snapshot available");
if (catalogSnapshotManager == null) {
throw new IllegalStateException("CatalogSnapshotManager not set");
}
return acquireReader(snapshot);
}

/**
* Acquires a dataFormatAwareReader on a specific catalog snapshot.
*/
private GatedCloseable<Reader> acquireReader(CatalogSnapshot catalogSnapshot) throws IOException {
catalogSnapshot.incRef();
GatedCloseable<CatalogSnapshot> snapshotRef = catalogSnapshotManager.acquireSnapshot();
try {
CatalogSnapshot catalogSnapshot = snapshotRef.get();
Map<DataFormat, Object> readers = new HashMap<>();
for (Map.Entry<DataFormat, EngineReaderManager<?>> entry : readerManagers.entrySet()) {
Object reader = entry.getValue().getReader(catalogSnapshot);
if (reader != null) {
readers.put(entry.getKey(), reader);
}
}
DataFormatAwareReader reader = new DataFormatAwareReader(catalogSnapshot, readers);
DataFormatAwareReader reader = new DataFormatAwareReader(catalogSnapshot, snapshotRef, readers);
return new GatedCloseable<>(reader, reader::close);
} catch (Exception e) {
catalogSnapshot.decRef();
snapshotRef.close();
throw e;
}
}
Expand All @@ -102,10 +97,16 @@ private GatedCloseable<Reader> acquireReader(CatalogSnapshot catalogSnapshot) th
@ExperimentalApi
public static class DataFormatAwareReader implements IndexReaderProvider.Reader {
private final CatalogSnapshot catalogSnapshot;
private final GatedCloseable<CatalogSnapshot> snapshotRef;
private final Map<DataFormat, Object> readers;

DataFormatAwareReader(CatalogSnapshot catalogSnapshot, Map<DataFormat, Object> readers) {
DataFormatAwareReader(
CatalogSnapshot catalogSnapshot,
GatedCloseable<CatalogSnapshot> snapshotRef,
Map<DataFormat, Object> readers
) {
this.catalogSnapshot = catalogSnapshot;
this.snapshotRef = snapshotRef;
this.readers = readers;
}

Expand All @@ -121,7 +122,11 @@ public CatalogSnapshot catalogSnapshot() {

@Override
public void close() {
catalogSnapshot.decRef();
try {
snapshotRef.close();
} catch (IOException e) {
throw new RuntimeException("Failed to release catalog snapshot reference", e);
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,8 @@
import org.opensearch.common.unit.TimeValue;
import org.opensearch.core.common.unit.ByteSizeValue;
import org.opensearch.index.VersionType;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.Indexer;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.mapper.DocumentMapperForType;
import org.opensearch.index.mapper.SourceToParse;
import org.opensearch.index.merge.MergeStats;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
import org.opensearch.common.logging.Loggers;
import org.opensearch.core.index.shard.ShardId;
import org.opensearch.index.engine.dataformat.MergeResult;
import org.opensearch.index.engine.exec.CatalogSnapshot;
import org.opensearch.index.engine.exec.Indexer;
import org.opensearch.index.engine.exec.Segment;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.util.ArrayDeque;
import java.util.Collection;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
package org.opensearch.index.engine.exec;

import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.io.IOException;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.io.IOException;
import java.util.Collection;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
package org.opensearch.index.engine.exec;

import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.io.IOException;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.index.engine.DataFormatAwareEngine;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.shard.ShardPath;

import java.io.IOException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.common.concurrent.GatedCloseable;
import org.opensearch.index.engine.dataformat.DataFormat;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;

import java.io.Closeable;
import java.io.IOException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import org.opensearch.index.engine.EngineException;
import org.opensearch.index.engine.LifecycleAware;
import org.opensearch.index.engine.SafeCommitInfo;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.translog.Translog;
import org.opensearch.index.translog.TranslogManager;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import org.opensearch.common.unit.TimeValue;
import org.opensearch.core.common.unit.ByteSizeValue;
import org.opensearch.index.engine.EngineException;
import org.opensearch.index.engine.exec.coord.CatalogSnapshot;
import org.opensearch.index.shard.ShardPath;

import java.io.IOException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,23 +9,58 @@
package org.opensearch.index.engine.exec;

import org.opensearch.common.annotation.ExperimentalApi;
import org.opensearch.core.common.io.stream.StreamInput;
import org.opensearch.core.common.io.stream.StreamOutput;
import org.opensearch.core.common.io.stream.Writeable;
import org.opensearch.index.engine.dataformat.DataFormat;

import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.function.Function;

/**
* Represents a segment in the catalog snapshot containing files grouped by data format.
* Each segment has a unique generation number and maintains searchable files organized by their data format type.
* This class is serializable and can be transmitted across nodes for replication and recovery operations.
*/
@ExperimentalApi
public record Segment(long generation, Map<String, WriterFileSet> dfGroupedSearchableFiles) {
public record Segment(long generation, Map<String, WriterFileSet> dfGroupedSearchableFiles) implements Writeable {

public Segment {
dfGroupedSearchableFiles = Map.copyOf(dfGroupedSearchableFiles);
}

/**
* Constructs a Segment by deserializing from a {@link StreamInput}.
*
* @param in the stream input to read from
* @param directoryResolver function that maps a data format name to its directory path
*/
public Segment(StreamInput in, Function<String, String> directoryResolver) throws IOException {
this(in.readLong(), readWriterFileSets(in, directoryResolver));
}

private static Map<String, WriterFileSet> readWriterFileSets(StreamInput in, Function<String, String> directoryResolver)
throws IOException {
int size = in.readVInt();
Map<String, WriterFileSet> map = new HashMap<>(size);
for (int i = 0; i < size; i++) {
String key = in.readString();
map.put(key, new WriterFileSet(in, directoryResolver.apply(key)));
}
return map;
}

@Override
public void writeTo(StreamOutput out) throws IOException {
out.writeLong(generation);
out.writeVInt(dfGroupedSearchableFiles.size());
for (Map.Entry<String, WriterFileSet> entry : dfGroupedSearchableFiles.entrySet()) {
out.writeString(entry.getKey());
entry.getValue().writeTo(out);
}
}

public static Builder builder(long generation) {
return new Builder(generation);
}
Expand Down
Loading
Loading