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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added
- Read `vortex.parquet.variant`, the encoding vortex-jni's default writer uses for Arrow `arrow.parquet.variant` columns: each row's Apache Variant `metadata`/`value` binaries come back as a struct ([#445](https://github.com/dfa1/vortex-java/issues/445)).
- Read `vortex.zstd_buffers`, Rust's opt-in buffer-level Zstd encoding (`zstd2026.02.0`) ([#444](https://github.com/dfa1/vortex-java/issues/444)).

### Fixed
- Null rows of Rust-written `vortex.varbin` columns read as empty values: the decoder ignored the validity child ([#444](https://github.com/dfa1/vortex-java/issues/444)).
- A corrupt `vortex.zstd` frame surfaced as the zstd binding's `ZstdException` instead of `VortexException` ([#444](https://github.com/dfa1/vortex-java/issues/444)).
- Dictionary-encoded Binary columns (e.g. the `value` child of `vortex.parquet.variant`) failed to read with a `ClassCastException` ([#445](https://github.com/dfa1/vortex-java/issues/445)).
- Null rows of Rust-written `vortex.list` columns read as empty lists: the decoder ignored the list's validity child.
- CSV export (`CsvExporter`, `vortex export`) failed on list-view columns with `unsupported array type for CSV export: ListViewArray`.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,8 @@
/// zone-map aggregates) has no members here. Encodings in no edition at all — `fastlanes.delta`,
/// `vortex.patched` — are writable only with the guard turned off, as in Rust.
///
/// vortex-java implements every `core`-family encoding (`vortex.parquet.variant` read only), but not
/// `zstd`'s `vortex.zstd_buffers`, which has no [EncodingId.WellKnown] constant yet and is named as
/// an [EncodingId.Custom] instead: the catalog mirrors upstream faithfully rather than being
/// truncated to what is implemented today.
/// vortex-java implements every encoding in the catalog; `vortex.parquet.variant` and
/// `vortex.zstd_buffers` are read only.
public final class Editions {

/// The baseline `core` edition: stable encodings writable by Vortex (Rust reference) 0.36.0.
Expand Down Expand Up @@ -85,10 +83,10 @@ public final class Editions {

/// The February 2026 draft edition of the `zstd` family, declared by Rust's `vortex-zstd` plugin:
/// buffer-level Zstd that keeps an array's buffer layout for GPU decompression. vortex-java
/// does not implement `vortex.zstd_buffers` yet, hence the [EncodingId.Custom].
/// reads `vortex.zstd_buffers` but does not write it.
public static final Edition ZSTD_2026_02_0 = new Edition(
new EditionId(EditionFamily.ZSTD, YearMonth.of(2026, 2), 0),
Set.of(new EncodingId.Custom("vortex.zstd_buffers")));
Set.of(EncodingId.VORTEX_ZSTD_BUFFERS));

/// Every declared edition, in Rust's declaration order (core declarations first, then the
/// plugin-declared families). Order matters:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,9 @@ enum WellKnown implements EncodingId {
/// Parquet Variant physical encoding (`vortex.parquet.variant`): per-row Apache Variant binary
/// `metadata` and `value` children, plus an optional shredded `typed_value` child.
VORTEX_PARQUET_VARIANT("vortex.parquet.variant"),
/// Buffer-level Zstd (`vortex.zstd_buffers`): each buffer of a wrapped array compressed
/// independently, keeping its layout; the `zstd` edition family, opt-in in Rust.
VORTEX_ZSTD_BUFFERS("vortex.zstd_buffers"),
;

// O(1) access to a WellKnown constant by its string representation
Expand Down Expand Up @@ -258,4 +261,6 @@ public String toString() {
WellKnown VORTEX_VARIANT = WellKnown.VORTEX_VARIANT;
/// Well-known `vortex.parquet.variant` id.
WellKnown VORTEX_PARQUET_VARIANT = WellKnown.VORTEX_PARQUET_VARIANT;
/// Well-known `vortex.zstd_buffers` id.
WellKnown VORTEX_ZSTD_BUFFERS = WellKnown.VORTEX_ZSTD_BUFFERS;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
package io.github.dfa1.vortex.core.proto;

import java.io.IOException;
import java.lang.foreign.MemorySegment;
import java.util.Arrays;
import javax.annotation.processing.Generated;

/// Generated from proto3 message {@code vortex.encodings.ZstdBuffersMetadata}.
/// Do not edit by hand — regenerate via {@code ./mvnw generate-sources -pl core -P regenerate-sources}.
/// @param inner_encoding_id field tag 1
/// @param inner_metadata field tag 2
/// @param uncompressed_sizes field tag 3
/// @param buffer_alignments field tag 4
/// @param child_dtypes field tag 5
/// @param child_lens field tag 6
@Generated("io.github.dfa1.vortex.protogen.CodeGen")
public record ProtoZstdBuffersMetadata(
String inner_encoding_id,
byte[] inner_metadata,
java.util.List<Long> uncompressed_sizes,
java.util.List<Integer> buffer_alignments,
java.util.List<ProtoDType> child_dtypes,
java.util.List<Long> child_lens
) {

/// Decodes a {@code vortex.encodings.ZstdBuffersMetadata} from a slice of a memory segment.
/// @param __seg backing segment
/// @param __off start offset in bytes
/// @param __len payload length in bytes
/// @return decoded record
/// @throws IOException if the slice is malformed or truncated
public static ProtoZstdBuffersMetadata decode(MemorySegment __seg, long __off, long __len) throws IOException {
ProtoReader r = new ProtoReader(__seg, __off, __len);
String inner_encoding_id = "";
byte[] inner_metadata = new byte[0];
java.util.List<Long> uncompressed_sizes = new java.util.ArrayList<>();
java.util.List<Integer> buffer_alignments = new java.util.ArrayList<>();
java.util.List<ProtoDType> child_dtypes = new java.util.ArrayList<>();
java.util.List<Long> child_lens = new java.util.ArrayList<>();
while (r.hasMore()) {
int tag = r.readVarint32();
switch (tag >>> 3) {
case 1 -> {
inner_encoding_id = r.readString();
}
case 2 -> {
inner_metadata = r.readBytes();
}
case 3 -> {
int wt = tag & 7;
if (wt == 2) {
int len = r.readVarint32();
java.util.List<Long> __target = uncompressed_sizes;
r.readPacked(len, reader -> __target.add(reader.readVarint64()));
} else {
uncompressed_sizes.add(r.readVarint64());
}
}
case 4 -> {
int wt = tag & 7;
if (wt == 2) {
int len = r.readVarint32();
java.util.List<Integer> __target = buffer_alignments;
r.readPacked(len, reader -> __target.add(reader.readVarint32()));
} else {
buffer_alignments.add(r.readVarint32());
}
}
case 5 -> {
MemorySegment __slice = r.readLenDelimSegment();
child_dtypes.add(ProtoDType.decode(__slice, 0, __slice.byteSize()));
}
case 6 -> {
int wt = tag & 7;
if (wt == 2) {
int len = r.readVarint32();
java.util.List<Long> __target = child_lens;
r.readPacked(len, reader -> __target.add(reader.readVarint64()));
} else {
child_lens.add(r.readVarint64());
}
}
default -> r.skipField(tag & 7);
}
}
return new ProtoZstdBuffersMetadata(inner_encoding_id, inner_metadata, uncompressed_sizes, buffer_alignments, child_dtypes, child_lens);
}

/// Encodes this record to a proto3-wire-format byte array.
/// @return encoded bytes
public byte[] encode() {
ProtoWriter w = new ProtoWriter();
encodeTo(w);
return w.toByteArray();
}

void encodeTo(ProtoWriter w) {
if (inner_encoding_id != null && !inner_encoding_id.isEmpty()) {
w.writeTag(1, 2);
w.writeString(inner_encoding_id);
}
if (inner_metadata != null && inner_metadata.length != 0) {
w.writeTag(2, 2);
w.writeBytes(inner_metadata);
}
if (!uncompressed_sizes.isEmpty()) {
w.writeTag(3, 2);
int __mark = w.beginLenDelim();
for (Long __v : uncompressed_sizes) {
w.writeVarint64(__v);
}
w.endLenDelim(__mark);
}
if (!buffer_alignments.isEmpty()) {
w.writeTag(4, 2);
int __mark = w.beginLenDelim();
for (Integer __v : buffer_alignments) {
w.writeVarint32(__v);
}
w.endLenDelim(__mark);
}
for (ProtoDType __v : child_dtypes) {
w.writeTag(5, 2);
int __mark = w.beginLenDelim();
__v.encodeTo(w);
w.endLenDelim(__mark);
}
if (!child_lens.isEmpty()) {
w.writeTag(6, 2);
int __mark = w.beginLenDelim();
for (Long __v : child_lens) {
w.writeVarint64(__v);
}
w.endLenDelim(__mark);
}
}

@Override
public boolean equals(Object __o) {
if (this == __o) {
return true;
}
if (!(__o instanceof ProtoZstdBuffersMetadata __that)) {
return false;
}
return java.util.Objects.equals(inner_encoding_id, __that.inner_encoding_id)
&& java.util.Arrays.equals(inner_metadata, __that.inner_metadata)
&& java.util.Objects.equals(uncompressed_sizes, __that.uncompressed_sizes)
&& java.util.Objects.equals(buffer_alignments, __that.buffer_alignments)
&& java.util.Objects.equals(child_dtypes, __that.child_dtypes)
&& java.util.Objects.equals(child_lens, __that.child_lens);
}

@Override
public int hashCode() {
int __h = 1;
__h = 31 * __h + java.util.Objects.hashCode(inner_encoding_id);
__h = 31 * __h + java.util.Arrays.hashCode(inner_metadata);
__h = 31 * __h + java.util.Objects.hashCode(uncompressed_sizes);
__h = 31 * __h + java.util.Objects.hashCode(buffer_alignments);
__h = 31 * __h + java.util.Objects.hashCode(child_dtypes);
__h = 31 * __h + java.util.Objects.hashCode(child_lens);
return __h;
}
}
9 changes: 9 additions & 0 deletions core/src/main/proto/encodings.proto
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,15 @@ message VariantMetadata {
optional vortex.dtype.DType shredded_dtype = 1;
}

message ZstdBuffersMetadata {
string inner_encoding_id = 1;
bytes inner_metadata = 2;
repeated uint64 uncompressed_sizes = 3;
repeated uint32 buffer_alignments = 4;
repeated vortex.dtype.DType child_dtypes = 5;
repeated uint64 child_lens = 6;
}

message ParquetVariantMetadata {
bool has_value = 1;
optional vortex.dtype.DType typed_value_dtype = 2;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -170,7 +170,7 @@ void owningEdition_zstdBuffers_isThePluginDeclaredZstdFamily() {
// Given — vortex.zstd_buffers joins Rust's zstd family, declared by the vortex-zstd
// plugin, not core (array-level vortex.zstd is core2025.06.0)
// When
Optional<Edition> result = Editions.owningEdition(new EncodingId.Custom("vortex.zstd_buffers"));
Optional<Edition> result = Editions.owningEdition(EncodingId.VORTEX_ZSTD_BUFFERS);

// Then
assertThat(result).contains(Editions.ZSTD_2026_02_0);
Expand Down
3 changes: 2 additions & 1 deletion docs/compatibility.md
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,7 @@ decimals ([#430](https://github.com/dfa1/vortex-java/pull/430)) and nulls in nul
| `vortex.patched` | `PatchedEncodingDecoder` | `PatchedEncodingEncoder` | ✅ | ✅ | Primitive PTypes; base + chunked patches (1024-elem blocks) |
| `vortex.variant` | `VariantEncodingDecoder` | `VariantEncodingEncoder` | ✅ | ✅ | Canonical container; constant / chunked-of-constants core + optional shredded child. Typed-scalar values only — nested objects need `parquet.variant` (ADR 0014) |
| `vortex.parquet.variant` | `ParquetVariantEncodingDecoder` | — | ✅ | ❌ | Read only: decodes to a struct `{metadata, value?, typed_value?}` of Apache Variant binaries (Arrow's `arrow.parquet.variant` storage shape), masked by row validity. Writing is not implemented; Java writes typed-scalar variants through `vortex.variant` (ADR 0014) |
| `vortex.zstd_buffers` | `ZstdBuffersEncodingDecoder` | — | ✅ | ❌ | Read only (`zstd2026.02.0`, opt-in in Rust): each buffer of any wrapped encoding is Zstd-decompressed at its declared alignment, then the inner encoding decodes with its own children. Needs the optional zstd binding, like `vortex.zstd`. Interop-tested against a Rust-written fixture (`scripts/fixtures/zstd-buffers`) |
| `vortex.onpair` | `OnPairEncodingDecoder` | `OnPairEncodingEncoder` | ✅ | ✅ | Utf8, Binary; `core2026.08.1`, a cascade candidate under the default edition (competes with FSST, as in Rust) |

### Decode shape
Expand Down Expand Up @@ -218,7 +219,7 @@ guarantee once frozen (ADR 0023) — a write-time/read-time policy, not part of
| `core2026.08.2` | `core` | `vortex.map` |
| `core2026.08.3` | `core` | `vortex.variant`, `vortex.parquet.variant` (read only) — **default write target** |
| `preview2026.08.0` | `preview` | nothing yet |
| `zstd2026.02.0` | `zstd` | `vortex.zstd_buffers` ❌ not implemented (buffer-level Zstd for GPU decode; declared by Rust's `vortex-zstd` plugin, opt-in) |
| `zstd2026.02.0` | `zstd` | `vortex.zstd_buffers` (read only; buffer-level Zstd for GPU decode; declared by Rust's `vortex-zstd` plugin, opt-in) |

Mirrors Rust's `vortex-edition` declarations at the pinned vortex-jni release (0.86.1);
`EditionCatalogParityIntegrationTest` fails the build on drift. `core` editions are frozen with a
Expand Down
5 changes: 2 additions & 3 deletions docs/reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -347,9 +347,8 @@ targeted edition is ever persisted into a `.vortex` file.
| `cumulativeMembers(Edition)` | The edition's own additions plus every earlier same-family edition's |
| `owningEdition(EncodingId)` | The edition an id first joined, or empty if it belongs to none |

vortex-java implements every `core`-family encoding (`vortex.parquet.variant` read only), but not
`zstd`'s `vortex.zstd_buffers`, which resolves to `EncodingId.Custom` and is stored in the catalog
anyway, mirroring upstream faithfully.
vortex-java implements every encoding in the catalog; `vortex.parquet.variant` and `zstd`'s
`vortex.zstd_buffers` are read only.
`fastlanes.delta` and `vortex.patched` belong to no edition, as in Rust.

### Writer integration (`WriteOptions#editions()`)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -222,14 +222,16 @@ private static List<String> readShredded(Path file) throws IOException {
while (iter.hasNext()) {
VariantArray column = iter.next().column("v");
StructArray struct = (StructArray) column.coreStorage();
VarBinArray value = (VarBinArray) struct.field(1);
MaskedArray value = (MaskedArray) struct.field(1);
MaskedArray typed = (MaskedArray) column.shredded();
for (long i = 0; i < struct.length(); i++) {
int row = rows.size();
// Exactly one of value / typed_value holds each row; the other is null.
assertThat(value.isValid(i)).as("value validity, row %d", row).isNotEqualTo(typed.isValid(i));
if (typed.isValid(i)) {
rows.add("typed:" + ((IntArray) typed.inner()).getInt(i));
} else {
assertThat(value.getBytes(i)).isEqualTo(int8OrString(row));
assertThat(((VarBinArray) value.inner()).getBytes(i)).isEqualTo(int8OrString(row));
rows.add("value:" + row);
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
package io.github.dfa1.vortex.integration;

import io.github.dfa1.vortex.inspect.InspectorTree;
import io.github.dfa1.vortex.reader.ReadRegistry;
import io.github.dfa1.vortex.reader.ScanOptions;
import io.github.dfa1.vortex.reader.Chunk;
import io.github.dfa1.vortex.reader.VortexReader;
import io.github.dfa1.vortex.reader.array.Array;
import io.github.dfa1.vortex.reader.array.IntArray;
import io.github.dfa1.vortex.reader.array.MaskedArray;
import io.github.dfa1.vortex.reader.array.VarBinArray;
import org.junit.jupiter.api.Test;

import java.io.IOException;
import java.net.URISyntaxException;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;

import static org.assertj.core.api.Assertions.assertThat;

/// Reads `fixtures/zstd_buffers.vortex`, written by Rust's own `ZstdBuffers::compress` under the
/// opt-in `zstd2026.02.0` edition (`scripts/fixtures/zstd-buffers`); vortex-jni cannot emit
/// `vortex.zstd_buffers`, so a checked-in fixture stands in for a jni-written file (#444).
///
/// Both columns wrap an inner array that has children as well as buffers (a primitive's validity,
/// a varbin's offsets and validity), so the decoder must hand the inner encoding its decompressed
/// buffers and its untouched children together.
class ZstdBuffersInteropIntegrationTest {

private static final int ROWS = 1_000;

@Test
void rustWritesZstdBuffers_javaReadsEveryRow() throws IOException, URISyntaxException {
// Given
Path file = Path.of(Objects.requireNonNull(
getClass().getResource("/fixtures/zstd_buffers.vortex")).toURI());

// When
List<String> result = new ArrayList<>(ROWS);
try (var reader = VortexReader.open(file, ReadRegistry.loadAll());
var iter = reader.scan(ScanOptions.columns("ints", "strs"))) {
while (iter.hasNext()) {
Chunk chunk = iter.next();
Array ints = chunk.column("ints");
Array strs = chunk.column("strs");
for (long i = 0; i < chunk.rowCount(); i++) {
result.add(render(ints, i, true) + "|" + render(strs, i, false));
}
}
}

// Then — the fixture's generator documents these values
try (var vf = VortexReader.open(file, ReadRegistry.loadAll())) {
assertThat(InspectorTree.build(vf).usedEncodings()).contains("vortex.zstd_buffers");
}
assertThat(result).hasSize(ROWS);
for (int i = 0; i < ROWS; i++) {
String ints = i % 7 == 0 ? "null" : String.valueOf(i * 3 - 500);
String strs = i % 5 == 0 ? "null" : "s" + i % 13;
assertThat(result.get(i)).as("row %d", i).isEqualTo(ints + "|" + strs);
}
}

private static String render(Array column, long i, boolean isInt) {
Array values = column;
if (column instanceof MaskedArray masked) {
if (!masked.isValid(i)) {
return "null";
}
values = masked.inner();
}
return isInt ? String.valueOf(((IntArray) values).getInt(i)) : ((VarBinArray) values).getString(i);
}
}
Binary file not shown.
Loading
Loading