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 @@ -16,6 +16,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Read and write `vortex.onpair` (unstable edition `unstable2026.06.0`): `tpch_orders.regular` and `clickbench_hits_5k.regular` from the v0.86.1 fixtures now scan, and `WriteOptions.withEdition(Editions.UNSTABLE_2026_06_0)` lets the cascade pick OnPair for string columns ([#425](https://github.com/dfa1/vortex-java/issues/425)).

### Changed
- With an unstable edition enabled (`WriteOptions.withEdition(Editions.UNSTABLE_2025_05_0)` or later), cascading writes can pick `fastlanes.delta`, its bases and deltas cascaded as in Rust; default writes never emit it ([#439](https://github.com/dfa1/vortex-java/pull/439)).
- The writer run-length encodes float columns too (`fastlanes.rle`, as Rust's float RLE scheme does), losslessly: `-0.0` and NaN payloads round-trip ([#438](https://github.com/dfa1/vortex-java/pull/438)).
- Cascading writes (`WriteOptions.cascading(n)`) compress run-end columns further: the run ends and values now go through the cascade (bit-packing, frame-of-reference, …) instead of being stored raw, as in Rust; about 1% smaller on the mixed-column size benchmark ([#435](https://github.com/dfa1/vortex-java/pull/435)).

Expand Down
2 changes: 1 addition & 1 deletion docs/compatibility.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ decimals ([#430](https://github.com/dfa1/vortex-java/pull/430)) and nulls in nul
| `vortex.datetimeparts` | `DateTimePartsEncodingDecoder` | `DateTimePartsEncodingEncoder` | ✅ | ✅ | |
| `vortex.pco` | `PcoEncodingDecoder` | `PcoEncodingEncoder` | ✅ | ✅ | Decode: all modes. Encode: Classic + Consecutive delta + IntMult; FloatMult/FloatQuant deferred |
| `fastlanes.bitpacked` | `BitpackedEncodingDecoder` | `BitpackedEncodingEncoder` | ✅ | ✅ | Unsigned integer PTypes |
| `fastlanes.delta` | `DeltaEncodingDecoder` | `DeltaEncodingEncoder` | ✅ | ✅ | Integer PTypes |
| `fastlanes.delta` | `DeltaEncodingDecoder` | `DeltaEncodingEncoder` | ✅ | ✅ | Integer PTypes. Unstable edition: cascading writes offer it (bases/deltas cascaded, as Rust's `DeltaScheme`) only with an unstable edition containing it |
| `fastlanes.for` | `FrameOfReferenceEncodingDecoder`| `FrameOfReferenceEncodingEncoder`| ✅ | ✅ | Integer PTypes |
| `fastlanes.rle` | `RleEncodingDecoder` | `RleEncodingEncoder` | ✅ | ✅ | Chunk-based RLE. Integers and floats (Rust's int and float RLE schemes); float runs compare raw bits, so -0.0 and NaN payloads round-trip. Cascades values/indices/offsets |
| `vortex.patched` | `PatchedEncodingDecoder` | `PatchedEncodingEncoder` | ✅ | ✅ | Primitive PTypes; base + chunked patches (1024-elem blocks) |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1315,6 +1315,43 @@ void javaWriter_rustReader_runEnd_i64(@TempDir Path tmp) throws IOException {
assertThat(decoded).containsExactly(data);
}

/// Cascaded delta (unstable edition, Rust's DeltaScheme, issue #410): bases and deltas go through
/// the cascade (FoR / bit-packing) instead of raw buffers. vortex-jni must read the column back on
/// a full scan, a mid-chunk row range, and a zone-pruned filter, which slice the delta children.
@Test
void javaWriter_jniReader_delta_cascading_unstableEdition(@TempDir Path tmp) throws IOException {
// Given — 2 equal chunks of ~1s jittered timestamps (delta wins over FoR here)
DType.Struct schema = new DType.Struct(List.of(ColumnName.of("id"), ColumnName.of("ts")),
List.of(DType.I64, DType.I64), false);
java.util.Random random = new java.util.Random(7);
int rows = 16_384;
long[] ts = new long[rows];
for (int i = 0; i < rows; i++) {
ts[i] = 1_700_000_000_000L + i * 1_000L + random.nextInt(1_000);
}
Path file = tmp.resolve("java_delta_cascade.vtx");
WriteOptions options = WriteOptions.cascading(3).withEdition(io.github.dfa1.vortex.core.model.Editions.UNSTABLE_2025_05_0);
try (var ch = FileChannel.open(file, StandardOpenOption.CREATE, StandardOpenOption.WRITE);
var sut = VortexWriter.create(ch, schema, options)) {
for (int start = 0; start < rows; start += 8_192) {
sut.writeChunk(Map.of(
ColumnName.of("id"), LongStream.range(start, start + 8_192L).toArray(),
ColumnName.of("ts"), Arrays.copyOfRange(ts, start, start + 8_192)));
}
}

// When
List<Object> full = readRowRange(file, "ts", 0, rows);
List<Object> range = readRowRange(file, "ts", 1_501, 12_345);
List<Object> filtered = readColumnFiltered(file, "ts", Expression.binary(Expression.BinaryOp.GTE,
Expression.column("id"), Expression.literal(12_000L)));

// Then
assertThat(full).containsExactlyInAnyOrderElementsOf(slice(Arrays.stream(ts).boxed(), 0, rows));
assertThat(range).containsExactlyInAnyOrderElementsOf(slice(Arrays.stream(ts).boxed(), 1_501, 12_345));
assertThat(filtered).containsExactlyInAnyOrderElementsOf(slice(Arrays.stream(ts).boxed(), 12_000, rows));
}

/// Float RLE (Rust's FloatRLEScheme): the writer used to refuse floats outright. Rust must read
/// the raw-bit runs back losslessly, -0.0 and a NaN payload included.
@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import io.github.dfa1.vortex.writer.encode.UuidExtensionEncoder;
import io.github.dfa1.vortex.writer.encode.DecimalBytePartsEncodingEncoder;
import io.github.dfa1.vortex.writer.encode.DecimalEncodingEncoder;
import io.github.dfa1.vortex.writer.encode.DeltaEncodingEncoder;
import io.github.dfa1.vortex.writer.encode.DictEncodingEncoder;
import io.github.dfa1.vortex.writer.encode.ZigZagEncodingEncoder;
import io.github.dfa1.vortex.writer.encode.ExtEncodingEncoder;
Expand Down Expand Up @@ -125,7 +126,7 @@ public final class VortexWriter implements Closeable {
private final Map<EncodingId, Integer> encodingIdx = new LinkedHashMap<>();
// Edition guard (issue #301): the cumulative member set of every WriteOptions#editions()
// family enabled for this writer, or empty when no edition is configured (guard off).
// editionExcluded is its complement over the concrete encoders this writer actually holds
// editionExcluded is its complement over every well-known encoding id plus the custom encoders this writer holds
// (encodings + cascadeCodecs) - seeded into every EncodeContext's initial `excluded` set so
// CascadingCompressor's existing per-candidate exclusion check (already consulted at every
// selection site, including nested competitions like a masked column's validity-bitmap
Expand Down Expand Up @@ -210,7 +211,7 @@ private static Set<EncodingId> editionAllowed(Map<EditionFamily, Edition> editio
return Set.copyOf(allowed);
}

/// The concrete encoder ids this writer actually holds that fall outside `allowed`. Seeding
/// Every encoder id — well-known, plus the custom ones this writer holds — outside `allowed`. Seeding
/// this into every [EncodeContext]'s initial exclusion set lets [CascadingCompressor]'s
/// existing per-candidate check skip them and fall back to the best remaining candidate,
/// instead of the edition guard only surfacing as a hard failure after encoding completes.
Expand All @@ -220,6 +221,15 @@ private static Set<EncodingId> editionExcluded(
return Set.of();
}
Set<EncodingId> excluded = new LinkedHashSet<>();
// Every well-known id, not only the encoders listed here: an encoder can carry its own private
// candidate list (Sparse's bool patch-index cascade offers fastlanes.delta), and those
// candidates must be filtered by the same edition policy.
for (EncodingId.WellKnown id : EncodingId.WellKnown.values()) {
if (!allowed.contains(id)) {
excluded.add(id);
}
}
// custom (non-well-known) encoders this writer holds
for (EncodingEncoder enc : encodings) {
if (!allowed.contains(enc.encodingId())) {
excluded.add(enc.encodingId());
Expand Down Expand Up @@ -274,6 +284,13 @@ private static List<EncodingEncoder> buildCascadeCodecs(WriteOptions options) {
codecs.add(new ZigZagEncodingEncoder());
codecs.add(new RunEndEncodingEncoder());
codecs.add(new RleEncodingEncoder());
// Delta (unstable edition, like Rust: "no edition includes fastlanes.delta yet, so the
// session's enabled editions decide") competes only when the caller opted into an unstable
// edition containing it; default writes never see it.
Edition unstable = options.editions().get(EditionFamily.UNSTABLE);
if (unstable != null && Editions.cumulativeMembers(unstable).contains(EncodingId.FASTLANES_DELTA)) {
codecs.add(new DeltaEncodingEncoder());
}
codecs.add(new SparseEncodingEncoder());
codecs.add(new DictEncodingEncoder());
codecs.add(new BitpackedEncodingEncoder());
Expand All @@ -288,7 +305,6 @@ private static List<EncodingEncoder> buildCascadeCodecs(WriteOptions options) {
// (e.g. taxi store_and_fwd_flag).
// OnPair (unstable edition) competes with FSST only when the caller opted into an unstable
// edition containing it; the default editions are core-only, so default writes never see it.
Edition unstable = options.editions().get(EditionFamily.UNSTABLE);
if (unstable != null && Editions.cumulativeMembers(unstable).contains(EncodingId.VORTEX_ONPAIR)) {
codecs.add(new OnPairEncodingEncoder());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import java.lang.foreign.MemorySegment;
import java.util.List;
import java.util.Set;

/// Write-only encoder for `fastlanes.delta`.
public final class DeltaEncodingEncoder implements EncodingEncoder {
Expand All @@ -32,8 +33,62 @@ public boolean accepts(DType dtype) {

@Override
public EncodeResult encode(DType dtype, Object data, EncodeContext ctx) {
PType ptype = ((DType.Primitive) dtype).ptype();
Deltas d = deltas(PrimitiveArrays.toLongs(data, ptype, EncodingId.FASTLANES_DELTA), ptype);
MemorySegment basesSeg = PrimitiveArrays.fromLongs(d.bases(), ptype, ctx.arena());
MemorySegment deltasSeg = PrimitiveArrays.fromLongs(d.deltas(), ptype, ctx.arena());
EncodeNode basesNode = EncodeNode.leaf(EncodingId.VORTEX_PRIMITIVE, 0);
EncodeNode deltasNode = EncodeNode.leaf(EncodingId.VORTEX_PRIMITIVE, 1);
EncodeNode root = new EncodeNode(EncodingId.FASTLANES_DELTA, MemorySegment.ofArray(d.metadata()),
new EncodeNode[]{basesNode, deltasNode}, new int[0]);
return new EncodeResult(root, List.of(EncodedBuffer.of(basesSeg, ptype), EncodedBuffer.of(deltasSeg, ptype)),
d.statsMin(), d.statsMax());
}

/// Barred from both children (issue #410, Rust's `DeltaScheme` descendant exclusions): delta
/// encoding data that is already delta encoded never pays off.
private static final Set<EncodingId> CHILDREN_EXCLUDED = Set.of(EncodingId.FASTLANES_DELTA);

/// Below one FastLanes chunk the transpose has nothing to work with (Rust's `MIN_DELTA_LEN`).
private static final int MIN_DELTA_LEN = FastLanes.CHUNK;

/// Cascading delta, mirroring Rust's `DeltaScheme`: the bases and deltas become open children
/// (Rust ids bases=0, deltas=1, also our wire order) that the compressor can FoR / bit-pack —
/// delta alone keeps the byte width. Same layout and metadata as [#encode].
@Override
public CascadeStep encodeCascade(DType dtype, Object data, EncodeContext ctx) {
PType ptype = ((DType.Primitive) dtype).ptype();
long[] longs = PrimitiveArrays.toLongs(data, ptype, EncodingId.FASTLANES_DELTA);
if (longs.length < MIN_DELTA_LEN) {
return CascadeStep.notApplicable();
}
Deltas d = deltas(longs, ptype);
EncodeNode partialRoot = new EncodeNode(EncodingId.FASTLANES_DELTA, MemorySegment.ofArray(d.metadata()),
new EncodeNode[]{null, null}, new int[0]);
DType childDtype = dtype.withNullable(false);
List<ChildSlot> slots = List.of(
new ChildSlot(childDtype, PrimitiveArrays.fromLongsArray(d.bases(), ptype, EncodingId.FASTLANES_DELTA),
0, CHILDREN_EXCLUDED),
new ChildSlot(childDtype, PrimitiveArrays.fromLongsArray(d.deltas(), ptype, EncodingId.FASTLANES_DELTA),
1, CHILDREN_EXCLUDED));
return new CascadeStep(partialRoot, List.of(), slots, d.statsMin(), d.statsMax(), true);
}

/// A column delta encoded in transposed FastLanes chunks.
///
/// @param bases per chunk, one base per lane
/// @param deltas per padded row, the wrapping difference from the previous row in its lane
/// @param paddedLen row count rounded up to a whole chunk
/// @param statsMin zone-map minimum, `null` when empty
/// @param statsMax zone-map maximum, `null` when empty
private record Deltas(long[] bases, long[] deltas, long paddedLen, byte[] statsMin, byte[] statsMax) {

byte[] metadata() {
return new ProtoDeltaMetadata(paddedLen, 0).encode();
}
}

private static Deltas deltas(long[] longs, PType ptype) {
int n = longs.length;
int typeBits = ptype.bits();
int lanes = FastLanes.lanes(ptype);
Expand Down Expand Up @@ -86,19 +141,9 @@ public EncodeResult encode(DType dtype, Object data, EncodeContext ctx) {
System.arraycopy(chunkDelta, 0, deltasAll, chunk * FastLanes.CHUNK, FastLanes.CHUNK);
}

MemorySegment basesSeg = PrimitiveArrays.fromLongs(basesAll, ptype, ctx.arena());
MemorySegment deltasSeg = PrimitiveArrays.fromLongs(deltasAll, ptype, ctx.arena());

byte[] metaBytes = new ProtoDeltaMetadata(paddedLen, 0).encode();

byte[] statsMin = n > 0 ? statsBytes(ptype, minVal) : null;
byte[] statsMax = n > 0 ? statsBytes(ptype, maxVal) : null;

EncodeNode basesNode = EncodeNode.leaf(EncodingId.VORTEX_PRIMITIVE, 0);
EncodeNode deltasNode = EncodeNode.leaf(EncodingId.VORTEX_PRIMITIVE, 1);
EncodeNode root = new EncodeNode(EncodingId.FASTLANES_DELTA, MemorySegment.ofArray(metaBytes),
new EncodeNode[]{basesNode, deltasNode}, new int[0]);
return new EncodeResult(root, List.of(EncodedBuffer.of(basesSeg, ptype), EncodedBuffer.of(deltasSeg, ptype)), statsMin, statsMax);
return new Deltas(basesAll, deltasAll, paddedLen, statsMin, statsMax);
}

private static void deltaChunk(long[] transposed, long[] bases, int lanes, int typeBits, long mask, long[] out) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,4 +85,83 @@ void withEdition_enablingUnstable_allowsTheForcedEncoder(@TempDir Path tmp) thro
assertThat(readAllLongs(vf, "ts")).containsExactly(data);
}
}

/// Delta joins the cascade like Rust's `DeltaScheme`: registered, but only offered when an
/// enabled edition contains `fastlanes.delta` (issue #410). Default writes must never emit it,
/// even on data it would win.
@Test
void cascading_defaultEditions_neverEmitDelta(@TempDir Path tmp) throws IOException {
// Given
Path file = tmp.resolve("delta_default.vtx");

// When
try (var ch = FileChannel.open(file, StandardOpenOption.CREATE, StandardOpenOption.WRITE);
var sut = VortexWriter.create(ch, I64_SCHEMA, WriteOptions.cascading(3))) {
sut.writeChunk(Map.of(ColumnName.of("ts"), jitteredTimestamps()));
}

// Then
try (var vf = VortexReader.open(file)) {
assertThat(vf.footer().arraySpecs()).doesNotContain(io.github.dfa1.vortex.core.model.EncodingId.FASTLANES_DELTA);
}
}

@Test
void cascading_unstableEdition_picksDeltaOnJitteredTimestamps(@TempDir Path tmp) throws IOException {
// Given — ~1s ticks with sub-second jitter: FoR needs the whole span (~23 bits), the
// transposed deltas only the jitter around the lane stride
Path file = tmp.resolve("delta_unstable_cascade.vtx");
long[] data = jitteredTimestamps();
WriteOptions options = WriteOptions.cascading(3).withEdition(Editions.UNSTABLE_2025_05_0);

// When
try (var ch = FileChannel.open(file, StandardOpenOption.CREATE, StandardOpenOption.WRITE);
var sut = VortexWriter.create(ch, I64_SCHEMA, options)) {
sut.writeChunk(Map.of(ColumnName.of("ts"), data));
}

// Then
try (var vf = VortexReader.open(file)) {
assertThat(vf.footer().arraySpecs()).contains(io.github.dfa1.vortex.core.model.EncodingId.FASTLANES_DELTA);
assertThat(readAllLongs(vf, "ts")).containsExactly(data);
}
}

private static long[] jitteredTimestamps() {
java.util.Random random = new java.util.Random(7);
long[] data = new long[8_192];
for (int i = 0; i < data.length; i++) {
data[i] = 1_700_000_000_000L + i * 1_000L + random.nextInt(1_000);
}
return data;
}

/// The edition policy must reach candidates an encoder keeps privately: Sparse compresses a
/// validity bitmap's patch indices over its own list, which offers `fastlanes.delta`. Before,
/// only the writer's own encoder lists were excluded, so once cascading delta could win there
/// a default-edition write failed with "fastlanes.delta: outside the configured edition(s)".
@Test
void cascading_defaultEditions_privateCandidateListsHonorTheEdition(@TempDir Path tmp) throws IOException {
// Given — a nullable column with a periodic null pattern: its validity goes sparse, and
// the regular patch indices are exactly what delta compresses best
DType.Struct schema = new DType.Struct(List.of(ColumnName.of("s")), List.of(new DType.Utf8(true)), false);
String[] categories = {"alpha", "beta", "gamma", "delta", "epsilon", "zeta", "eta", "theta"};
java.util.Random random = new java.util.Random(42);
String[] data = new String[50_000];
for (int i = 0; i < data.length; i++) {
data[i] = i % 10 == 0 ? null : categories[random.nextInt(categories.length)];
}
Path file = tmp.resolve("sparse_validity.vtx");

// When
try (var ch = FileChannel.open(file, StandardOpenOption.CREATE, StandardOpenOption.WRITE);
var sut = VortexWriter.create(ch, schema, WriteOptions.cascading(3).withGlobalDict(false))) {
sut.writeChunk(Map.of(ColumnName.of("s"), data));
}

// Then
try (var vf = VortexReader.open(file)) {
assertThat(vf.footer().arraySpecs()).doesNotContain(io.github.dfa1.vortex.core.model.EncodingId.FASTLANES_DELTA);
}
}
}
Loading
Loading