Skip to content

Commit e71563a

Browse files
dfa1claude
andcommitted
fix(writer): null as a global-dict entry for string and binary columns
## fix(writer): store null as a global-dict entry for string and binary columns Same shape as the numeric path: a nullable Utf8/Binary global dict now writes null as one more pool entry, invalid in the pool's validity, with non-nullable U16 codes and is_nullable_codes = false, as Rust's dict layout writer does. The pool (FSST/VarBin/Zstd competition) is now masked when it holds the null entry, so writeSegment's masked override honors the cascade exclusions too: the pool must never become a dict itself, which the reader cannot unwrap. Taxi 2024-01: 42,342,166 -> 42,340,950 bytes (vortex-jni 44,463,892) store_and_fwd_flag 66,460 -> 64,990 (jni 47,536: 1.40x -> 1.37x) pool and layout metadata now match vortex-jni's shape; the rest of the gap is 46 code chunks of 64k rows vs Rust's 6 of 512k OHLC 10M rows: unchanged at 58,230,658 bytes. JavaVsJniWriteBenchmark.javaWriteCascading, 5 forks: 0.769 +- 0.026 -> 0.767 +- 0.017 ops/s (per fork 0.769 0.768 0.805 0.742 0.751) gc.alloc.rate.norm 2.50 -> 2.50 GB/op Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015R8cf8YM8XFQqrkyhSAtYr
1 parent 7ad1ba9 commit e71563a

3 files changed

Lines changed: 37 additions & 17 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
88
## [Unreleased]
99

1010
### Changed
11-
- Nullable low-cardinality numeric columns written with a global dictionary store null as one more dictionary entry, invalid in the values' validity, with non-nullable codes, as Rust's dict layout does: no row mask on the codes, and their encoding is chosen on the real codes rather than placeholder ones. Global-dict codes are always U16, Rust's default, whatever the dictionary size ([#467](https://github.com/dfa1/vortex-java/pull/467)).
11+
- Nullable low-cardinality numeric, string and binary columns written with a global dictionary store null as one more dictionary entry, invalid in the values' validity, with non-nullable codes, as Rust's dict layout does: no row mask on the codes, and their encoding is chosen on the real codes rather than placeholder ones. Global-dict codes are always U16, Rust's default, whatever the dictionary size ([#467](https://github.com/dfa1/vortex-java/pull/467), [#468](https://github.com/dfa1/vortex-java/pull/468)).
1212
- Cascading writes are about 35% faster and allocate half as much (10M-row OHLC: 0.44 → 0.60 writes/s), to the same bytes: as in Rust, cascade children are no longer trial-encoded with candidates their stats rule out, and RLE/RunEnd skip arrays whose runs average under 4 values ([#465](https://github.com/dfa1/vortex-java/pull/465)).
1313
- **Breaking:** `EncodingEncoder#expectedRatio` takes `(DType, ArrayAndStats, EncodeContext)` instead of `(DType, Object, ArrayStats)`, as Rust's `expected_compression_ratio` takes `ArrayAndStats` and the compressor context: the cascade's stats scan runs on the first `stats()` call and is shared by every candidate, and `EncodeContext#sample()`/`finishedCascading()` tell a verdict where it is. `Estimate` is a sealed interface whose `Estimate.ratio(double)` competes with sampled candidates without encoding, as Rust's `EstimateVerdict::Ratio` does ([#464](https://github.com/dfa1/vortex-java/pull/464), [#466](https://github.com/dfa1/vortex-java/pull/466)).
1414
- **Breaking:** default writes emit Rust's `vortex.zoned` zone map instead of the legacy `vortex.stats`: one zone per 8192 rows regardless of `writeChunk` batches, with Rust's per-dtype stats (64-byte string bounds, `nan_count`) and no zone sum. Calcite and `ZoneReducer` aggregate push-down needs the legacy layout, so on these files it falls back to a scan; target `Editions.CORE_2025_10_0` to keep it. Read new files with this release or later ([#447](https://github.com/dfa1/vortex-java/issues/447)).

‎writer/src/main/java/io/github/dfa1/vortex/writer/DictColumnState.java‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -452,8 +452,7 @@ Object buildFrequencyRankedUniqueArray(int[] remap) {
452452
/// Emits one chunk's wire codes ([#CODES_PTYPE]) from its buffered `short[]` first-seen codes,
453453
/// optionally translated through a frequency remap (primitive path). A null slot (validity[i]
454454
/// false) emits `nullCode` unconditionally, never remapped: the primitive path points it at the
455-
/// pool's invalid null entry, as Rust's dict builder does; the Utf8/Binary path passes `0`, which
456-
/// the reader ignores because its codes child is masked by the same validity.
455+
/// pool's invalid null entry, as Rust's dict builder does.
457456
///
458457
/// @param buffered the buffered first-seen codes for one chunk
459458
/// @param remap the first-seen -> frequency-rank remap, or `null` to emit codes unchanged
@@ -474,9 +473,10 @@ static short[] emitCodes(short[] buffered, int[] remap, boolean[] validity, int
474473
}
475474

476475
/// The values pool with one more slot, for the invalid null entry Rust's dict builder adds
477-
/// when it meets a null: the slot holds a zero placeholder and is masked off by the caller.
476+
/// when it meets a null: the slot holds a zero (or `null`) placeholder and is masked off by the
477+
/// caller.
478478
///
479-
/// @param values the frequency-ranked distinct values (a primitive array)
479+
/// @param values the distinct values (a primitive array, `String[]` or `byte[][]`)
480480
/// @return a copy one element longer
481481
static Object withNullSlot(Object values) {
482482
return switch (values) {
@@ -486,7 +486,10 @@ static Object withNullSlot(Object values) {
486486
case byte[] a -> Arrays.copyOf(a, a.length + 1);
487487
case double[] a -> Arrays.copyOf(a, a.length + 1);
488488
case float[] a -> Arrays.copyOf(a, a.length + 1);
489-
default -> throw new IllegalStateException("not a primitive values pool: " + values.getClass());
489+
// Utf8/Binary: the null slot holds a real null, as ChunkImpl leaves it at invalid rows
490+
case String[] a -> Arrays.copyOf(a, a.length + 1);
491+
case byte[][] a -> Arrays.copyOf(a, a.length + 1);
492+
default -> throw new IllegalStateException("not a values pool: " + values.getClass());
490493
};
491494
}
492495

‎writer/src/main/java/io/github/dfa1/vortex/writer/VortexWriter.java‎

Lines changed: 28 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -691,6 +691,9 @@ private int writeSegment(DType dtype, Object data, EncodingEncoder encodingOverr
691691
EncodeContext encodeCtx = options.allowedCascading() > 0
692692
? EncodeContext.ofDepth(options.allowedCascading(), arena, cascadeRegistry, editionExcluded)
693693
: EncodeContext.of(arena, defaultRegistry, editionExcluded);
694+
// A wrapping override (masked) cascades its inner values, so the exclusions hold
695+
// there too -- a nullable global-dict pool must never become a dict itself.
696+
encodeCtx = encodeCtx.withExcluded(excludedFromCascade);
694697
result = encodingOverride.encode(dtype, data, encodeCtx);
695698
} else if (options.allowedCascading() > 0) {
696699
EncodeContext encodeCtx = EncodeContext.ofDepth(options.allowedCascading(), arena, cascadeRegistry, editionExcluded);
@@ -1366,31 +1369,45 @@ private void writeGlobalDictColumn(ColumnName colName, DictColumnState state) th
13661369
}
13671370

13681371
private void writeGlobalDictVarBinColumn(ColumnName colName, DictColumnState state) throws IOException {
1369-
int dictSize = state.cardinality();
1372+
// Null is one more pool entry, invalid, with non-nullable codes -- Rust's dict layout
1373+
// shape, as for numeric columns (writeGlobalDictColumn).
1374+
boolean hasNulls = state.chunkNullCounts().stream().anyMatch(c -> c > 0);
1375+
int distinct = state.cardinality();
1376+
int dictSize = hasNulls ? distinct + 1 : distinct;
13701377

13711378
// Utf8/Binary assigns codes in first-seen order with no frequency sort, so the incremental map's
13721379
// order already matches — no remap pass (ADR 0021). Compress the distinct-values pool
13731380
// through the normal Utf8/Binary competition (FSST/VarBin/Zstd) so it captures substring
13741381
// redundancy across dictionary entries (#299), but exclude Dict so the cascade never wraps
13751382
// the (all-unique-by-construction) dictionary in another dict the reader cannot unwrap. At
1376-
// cascade depth 0 there is no competition to run, so force flat VarBin as before.
1383+
// cascade depth 0 there is no competition to run, so force flat VarBin as before; a pool
1384+
// with a null entry is masked there instead, its values still VarBin.
13771385
Object uniques = state.varBinUniques();
1378-
int valuesSegIdx = options.allowedCascading() > 0
1379-
? writeSegment(state.dtype(), uniques, null, Set.of(EncodingId.VORTEX_DICT))
1380-
: writeSegment(state.dtype(), uniques, new VarBinEncodingEncoder());
1386+
Object pool = uniques;
1387+
if (hasNulls) {
1388+
boolean[] poolValidity = new boolean[dictSize];
1389+
Arrays.fill(poolValidity, 0, distinct, true);
1390+
pool = new NullableData(DictColumnState.withNullSlot(uniques), poolValidity);
1391+
}
1392+
int valuesSegIdx;
1393+
if (options.allowedCascading() > 0) {
1394+
valuesSegIdx = writeSegment(state.dtype(), pool, null, Set.of(EncodingId.VORTEX_DICT));
1395+
} else if (hasNulls) {
1396+
valuesSegIdx = writeSegment(state.dtype(), pool, new MaskedEncodingEncoder(), Set.of(EncodingId.VORTEX_DICT));
1397+
} else {
1398+
valuesSegIdx = writeSegment(state.dtype(), pool, new VarBinEncodingEncoder());
1399+
}
13811400

1382-
DType codesDtype = new DType.Primitive(DictColumnState.CODES_PTYPE, state.nullable());
1401+
DType codesDtype = new DType.Primitive(DictColumnState.CODES_PTYPE, false);
13831402
List<Integer> codesSegIdxes = new ArrayList<>();
13841403
for (int c = 0; c < state.chunkCount(); c++) {
1385-
boolean[] validity = state.chunkValidity(c);
1386-
Object codesArr = DictColumnState.emitCodes(state.chunkCodes(c), null, validity, 0);
1387-
Object codesData = validity != null ? new NullableData(codesArr, validity) : codesArr;
1388-
codesSegIdxes.add(writeSegment(codesDtype, codesData));
1404+
Object codesArr = DictColumnState.emitCodes(state.chunkCodes(c), null, state.chunkValidity(c), distinct);
1405+
codesSegIdxes.add(writeSegment(codesDtype, codesArr));
13891406
}
13901407

13911408
dictColRefs.put(colName, new DictColRef(valuesSegIdx, dictSize, codesSegIdxes,
13921409
state.chunkRowCounts(), state.chunkNullCounts(),
1393-
state.chunkStatsMin(), state.chunkStatsMax(), state.chunkStatsSum(), state.nullable()));
1410+
state.chunkStatsMin(), state.chunkStatsMax(), state.chunkStatsSum(), false));
13941411
}
13951412

13961413
private record SegRef(long offset, long len) {

0 commit comments

Comments
 (0)