Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -274,4 +274,8 @@ public double getDoubleLE(int index) {
}
return super.getDoubleLE(index);
}

public BytesReference[] references() {
return references;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,10 @@ public ReleasableBytesReference retainedSlice(int from, int length) {
return new ReleasableBytesReference(slice, refCounted);
}

public BytesReference delegate() {
return delegate;
}

@Override
public void close() {
refCounted.decRef();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,15 @@

import org.elasticsearch.TransportVersion;
import org.elasticsearch.common.bytes.BytesReference;
import org.elasticsearch.common.bytes.CompositeBytesReference;
import org.elasticsearch.common.bytes.ReleasableBytesReference;
import org.elasticsearch.common.util.PageCacheRecycler;
import org.elasticsearch.core.Nullable;
import org.elasticsearch.core.Releasable;

import java.io.IOException;
import java.io.UncheckedIOException;
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;

Expand Down Expand Up @@ -203,8 +206,7 @@ public boolean isSerialized() {

@Override
public long getSerializedSize() {
// We're already serialized
return serialized.length();
return pageAlignedRamUsedByReferenceBytes(serialized);
}

@Override
Expand All @@ -213,6 +215,22 @@ public void close() {
}
}

/**
* Over-estimates retained RAM by rounding each component up to a
* {@link PageCacheRecycler#BYTE_PAGE_SIZE} page (composites sum per component).
* Matches how Netty retains page-sized buffers rather than exact payload lengths.
*/
static long pageAlignedRamUsedByReferenceBytes(BytesReference bytes) {
if (bytes instanceof ReleasableBytesReference r) {
return pageAlignedRamUsedByReferenceBytes(r.delegate());
}
if (bytes instanceof CompositeBytesReference composited) {
return Arrays.stream(composited.references()).mapToLong(DelayableWriteable::pageAlignedRamUsedByReferenceBytes).sum();
}
final long numPages = (bytes.length() + PageCacheRecycler.BYTE_PAGE_SIZE - 1L) / PageCacheRecycler.BYTE_PAGE_SIZE;
return numPages * PageCacheRecycler.BYTE_PAGE_SIZE;
}

/**
* Returns the serialized size in bytes of the provided {@link Writeable}.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,11 @@
package org.elasticsearch.common.io.stream;

import org.elasticsearch.TransportVersion;
import org.elasticsearch.common.bytes.BytesArray;
import org.elasticsearch.common.bytes.BytesReference;
import org.elasticsearch.common.bytes.CompositeBytesReference;
import org.elasticsearch.common.bytes.ReleasableBytesReference;
import org.elasticsearch.common.util.PageCacheRecycler;
import org.elasticsearch.test.ESTestCase;
import org.elasticsearch.test.TransportVersionUtils;

Expand Down Expand Up @@ -124,9 +129,12 @@ public void testRoundTripFromReferencingWithNamedWriteable() throws IOException
}

public void testRoundTripFromDelayed() throws IOException {
Example e = new Example(randomAlphaOfLength(5));
Example e = new Example(randomAlphaOfLengthBetween(100, 1000));
DelayableWriteable<Example> original = DelayableWriteable.referencing(e).asSerialized(Example::new, writableRegistry());
assertTrue(original.isSerialized());
long length = DelayableWriteable.getSerializedSize(e);
long page = PageCacheRecycler.BYTE_PAGE_SIZE;
assertThat(original.getSerializedSize(), equalTo(((length + page - 1) / page) * page));
roundTripTestCase(original, Example::new);
}

Expand Down Expand Up @@ -165,6 +173,26 @@ public void testAsSerializedIsNoopOnSerialized() throws IOException {
assertSame(d, d.asSerialized(Example::new, writableRegistry()));
}

public void testPageAlignedRamUsedByReferenceBytes() {
final int page = PageCacheRecycler.BYTE_PAGE_SIZE;
assertThat(DelayableWriteable.pageAlignedRamUsedByReferenceBytes(BytesArray.EMPTY), equalTo(0L));
assertThat(DelayableWriteable.pageAlignedRamUsedByReferenceBytes(new BytesArray(new byte[1])), equalTo((long) page));
assertThat(DelayableWriteable.pageAlignedRamUsedByReferenceBytes(new BytesArray(new byte[page])), equalTo((long) page));
assertThat(DelayableWriteable.pageAlignedRamUsedByReferenceBytes(new BytesArray(new byte[page + 1])), equalTo(2L * page));

assertThat(
DelayableWriteable.pageAlignedRamUsedByReferenceBytes(ReleasableBytesReference.wrap(new BytesArray(new byte[1]))),
equalTo((long) page)
);

BytesReference composite = CompositeBytesReference.of(
new BytesArray(new byte[1]),
new BytesArray(new byte[page / 2]),
new BytesArray(new byte[page + 1])
);
assertThat(DelayableWriteable.pageAlignedRamUsedByReferenceBytes(composite), equalTo(4L * page));
}

private <T extends Writeable> void roundTripTestCase(DelayableWriteable<T> original, Writeable.Reader<T> reader) throws IOException {
DelayableWriteable<T> roundTripped = roundTrip(original, reader, TransportVersion.current());
assertThat(roundTripped.expand(), equalTo(original.expand()));
Expand Down