Add columnar V3 payload encoding tags

This commit is contained in:
Paul Masurel
2026-09-15 12:37:35 +02:00
parent f5069e5aae
commit 0bc563255b
12 changed files with 318 additions and 22 deletions
Binary file not shown.
Binary file not shown.
Binary file not shown.
+64 -1
View File
@@ -14,7 +14,7 @@ use crate::column_values::{
load_u64_based_column_values, serialize_column_values_u128, serialize_u64_based_column_values,
};
use crate::iterable::Iterable;
use crate::{StrColumn, Version};
use crate::{PayloadEncoding, StrColumn, Version};
pub fn serialize_column_mappable_to_u128<T: MonotonicallyMappableToU128>(
column_index: SerializableColumnIndex<'_>,
@@ -106,8 +106,43 @@ pub fn open_column_u128_as_compact_u64(
}
pub fn open_column_bytes(data: OwnedBytes, format_version: Version) -> io::Result<BytesColumn> {
let data = match format_version {
Version::V1 | Version::V2 => data,
Version::V3 => {
if data.is_empty() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"missing string/byte payload encoding tag",
));
}
let (encoding_bytes, payload) = data.split(1);
let encoding = PayloadEncoding::try_from_code(encoding_bytes.as_slice()[0])
.map_err(io::Error::from)?;
match encoding {
PayloadEncoding::Dictionary => payload,
PayloadEncoding::Plain => {
return Err(io::Error::new(
io::ErrorKind::Unsupported,
"plain string/byte column decoding is not implemented yet",
));
}
}
}
};
if data.len() < 4 {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"truncated dictionary string/byte column payload",
));
}
let (body, dictionary_len_bytes) = data.rsplit(4);
let dictionary_len = u32::from_le_bytes(dictionary_len_bytes.as_slice().try_into().unwrap());
if dictionary_len as usize > body.len() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"dictionary length exceeds string/byte column payload",
));
}
let (dictionary_bytes, column_bytes) = body.split(dictionary_len as usize);
let dictionary = Arc::new(Dictionary::from_bytes(dictionary_bytes)?);
let term_ord_column = crate::column::open_column_u64::<u64>(column_bytes, format_version)?;
@@ -125,3 +160,31 @@ pub fn open_column_str(data: OwnedBytes, format_version: Version) -> io::Result<
};
Ok(DictionaryEncodedStrColumn::wrap(bytes_column).into())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_v3_payload_encoding_tag_errors() {
let error = open_column_bytes(OwnedBytes::new(Vec::new()), Version::V3).unwrap_err();
assert_eq!(error.kind(), io::ErrorKind::InvalidData);
let error = open_column_bytes(OwnedBytes::new(vec![u8::MAX]), Version::V3).unwrap_err();
assert_eq!(error.kind(), io::ErrorKind::InvalidData);
let error = open_column_bytes(
OwnedBytes::new(vec![PayloadEncoding::Plain.to_code()]),
Version::V3,
)
.unwrap_err();
assert_eq!(error.kind(), io::ErrorKind::Unsupported);
let error = open_column_bytes(
OwnedBytes::new(vec![PayloadEncoding::Dictionary.to_code()]),
Version::V3,
)
.unwrap_err();
assert_eq!(error.kind(), io::ErrorKind::InvalidData);
}
}
@@ -55,7 +55,7 @@ pub fn open_multivalued_index(
start_index_column,
}))
}
Version::V2 => {
Version::V2 | Version::V3 => {
let (body_bytes, optional_index_len) = bytes.rsplit(4);
let optional_index_len =
u32::from_le_bytes(optional_index_len.as_slice().try_into().unwrap());
+6 -3
View File
@@ -23,13 +23,14 @@ pub fn parse_footer(footer_bytes: [u8; VERSION_FOOTER_NUM_BYTES]) -> Result<Vers
Version::try_from_bytes(footer_bytes[0..4].try_into().unwrap())
}
pub const CURRENT_VERSION: Version = Version::V2;
pub const CURRENT_VERSION: Version = Version::V3;
#[derive(Debug, Copy, Clone, Eq, PartialEq)]
#[repr(u32)]
pub enum Version {
V1 = 1u32,
V2 = 2u32,
V3 = 3u32,
}
impl Display for Version {
@@ -37,6 +38,7 @@ impl Display for Version {
match self {
Version::V1 => write!(f, "v1"),
Version::V2 => write!(f, "v2"),
Version::V3 => write!(f, "v3"),
}
}
}
@@ -51,6 +53,7 @@ impl Version {
match code {
1u32 => Ok(Version::V1),
2u32 => Ok(Version::V2),
3u32 => Ok(Version::V3),
_ => Err(InvalidData),
}
}
@@ -65,7 +68,7 @@ mod tests {
#[test]
fn test_footer_deserialization() {
let parsed_version: Version = parse_footer(footer()).unwrap();
assert_eq!(Version::V2, parsed_version);
assert_eq!(Version::V3, parsed_version);
}
#[test]
@@ -83,6 +86,6 @@ mod tests {
valid_versions.insert(i);
}
}
assert_eq!(valid_versions.len(), 2);
assert_eq!(valid_versions.len(), 3);
}
}
@@ -7,9 +7,11 @@ use super::term_merger::{TermMerger, TermsWithSegmentOrd};
use crate::column::serialize_column_mappable_to_u64;
use crate::column_index::SerializableColumnIndex;
use crate::iterable::Iterable;
use crate::{BytesColumn, DictionaryEncodedBytesColumn, MergeRowOrder, ShuffleMergeOrder};
use crate::{
BytesColumn, DictionaryEncodedBytesColumn, MergeRowOrder, PayloadEncoding, ShuffleMergeOrder,
};
// Serialize [Dictionary, Column, dictionary num bytes U32::LE]
// V3 serialize [PayloadEncoding, Dictionary, Column, dictionary num bytes U32::LE]
// Column: [Column Index, Column Values, column index num bytes U32::LE]
pub fn merge_bytes_or_str_column(
column_index: SerializableColumnIndex<'_>,
@@ -17,7 +19,9 @@ pub fn merge_bytes_or_str_column(
merge_row_order: &MergeRowOrder,
output: &mut impl Write,
) -> io::Result<()> {
// Serialize dict and generate mapping for values
output.write_all(&[PayloadEncoding::Dictionary.to_code()])?;
// Serialize dict and generate mapping for values.
// The encoding tag is intentionally excluded from `dictionary_num_bytes`.
let mut output = CountingWriter::wrap(output);
// TODO !!! Remove useless terms.
let term_ord_mapping = serialize_merged_dict(bytes_columns, merge_row_order, &mut output)?;
+4 -2
View File
@@ -22,7 +22,7 @@ use crate::columnar::writer::column_writers::{
use crate::columnar::writer::value_index::{IndexBuilder, PreallocatedIndexBuilders};
use crate::dictionary::{DictionaryBuilder, TermIdMapping, UnorderedId};
use crate::value::{Coerce, NumericalType, NumericalValue};
use crate::{Cardinality, RowId};
use crate::{Cardinality, PayloadEncoding, RowId};
/// This is a set of buffers that are used to temporarily write the values into before passing them
/// to the fast field codecs.
@@ -532,7 +532,7 @@ fn collect_sort_order_from_ops<V, K: Clone>(
.collect()
}
// Serialize [Dictionary, Column, dictionary num bytes U32::LE]
// V3 serialize [PayloadEncoding, Dictionary, Column, dictionary num bytes U32::LE]
// Column: [Column Index, Column Values, column index num bytes U32::LE]
#[expect(clippy::too_many_arguments)]
fn serialize_bytes_or_str_column(
@@ -550,6 +550,8 @@ fn serialize_bytes_or_str_column(
u64_values,
..
} = buffers;
let mut wrt = wrt;
wrt.write_all(&[PayloadEncoding::Dictionary.to_code()])?;
let mut counting_writer = CountingWriter::wrap(wrt);
let term_id_mapping: TermIdMapping =
dictionary_builder.serialize(arena, &mut counting_writer)?;
+154 -2
View File
@@ -3,11 +3,12 @@ use std::path::PathBuf;
use itertools::Itertools;
use crate::{
CURRENT_VERSION, Cardinality, Column, ColumnarReader, DynamicColumn, StackMergeOrder,
merge_columnar,
CURRENT_VERSION, Cardinality, Column, ColumnarReader, DictionaryEncodedBytesColumn,
DictionaryEncodedStrColumn, DynamicColumn, PayloadEncoding, StackMergeOrder, merge_columnar,
};
const NUM_DOCS: u32 = u16::MAX as u32;
const STRING_BYTES_NUM_DOCS: u32 = 4;
fn generate_columnar(num_docs: u32, value_offset: u64) -> Vec<u8> {
use crate::ColumnarWriter;
@@ -32,6 +33,32 @@ fn generate_columnar(num_docs: u32, value_offset: u64) -> Vec<u8> {
wrt
}
fn generate_string_bytes_columnar() -> Vec<u8> {
use crate::ColumnarWriter;
let mut columnar_writer = ColumnarWriter::default();
for doc in 0..STRING_BYTES_NUM_DOCS {
columnar_writer.record_str(doc, "str_full", &format!("str-full-{doc}"));
columnar_writer.record_bytes(doc, "bytes_full", &[doc as u8, 0, 255]);
if doc.is_multiple_of(2) {
columnar_writer.record_str(doc, "str_optional", &format!("str-optional-{doc}"));
columnar_writer.record_bytes(doc, "bytes_optional", &[doc as u8, 1, 254]);
}
columnar_writer.record_str(doc, "str_multi", &format!("str-multi-{doc}-a"));
columnar_writer.record_str(doc, "str_multi", &format!("str-multi-{doc}-b"));
columnar_writer.record_bytes(doc, "bytes_multi", &[doc as u8, 2, 0]);
columnar_writer.record_bytes(doc, "bytes_multi", &[doc as u8, 2, 255]);
}
let mut output = Vec::new();
columnar_writer
.serialize(STRING_BYTES_NUM_DOCS, None, &mut output)
.unwrap();
output
}
#[test]
/// Writes a columnar for the CURRENT_VERSION to disk.
fn create_format() {
@@ -63,6 +90,131 @@ fn test_format_v2() {
test_format(&path);
}
#[test]
fn test_format_v3() {
let path = path_for_version("v3");
test_format(&path);
}
#[test]
fn test_string_bytes_format_v1() {
test_string_bytes_format("v1");
}
#[test]
fn test_string_bytes_format_v2() {
test_string_bytes_format("v2");
}
fn test_string_bytes_format(version: &str) {
let fixture_path = format!("./compat_tests_data/{version}_string_bytes.columnar");
let fixture_reader = ColumnarReader::open(std::fs::read(fixture_path).unwrap()).unwrap();
check_string_bytes_columns(&fixture_reader, 1);
let current_reader = ColumnarReader::open(generate_string_bytes_columnar()).unwrap();
check_string_bytes_columns(&current_reader, 1);
let readers = [&fixture_reader, &current_reader];
let merge_row_order = StackMergeOrder::stack(&readers);
let mut output = Vec::new();
merge_columnar(&readers, &[], merge_row_order.into(), &mut output).unwrap();
let merged_reader = ColumnarReader::open(output).unwrap();
check_string_bytes_columns(&merged_reader, 2);
}
fn check_string_bytes_columns(reader: &ColumnarReader, repetitions: u32) {
let num_docs = STRING_BYTES_NUM_DOCS * repetitions;
let str_full = open_str_column(reader, "str_full");
assert_eq!(str_full.get_cardinality(), Cardinality::Full);
let str_optional = open_str_column(reader, "str_optional");
assert_eq!(str_optional.get_cardinality(), Cardinality::Optional);
let str_multi = open_str_column(reader, "str_multi");
assert_eq!(str_multi.get_cardinality(), Cardinality::Multivalued);
let bytes_full = open_bytes_column(reader, "bytes_full");
assert_eq!(bytes_full.get_cardinality(), Cardinality::Full);
let bytes_optional = open_bytes_column(reader, "bytes_optional");
assert_eq!(bytes_optional.get_cardinality(), Cardinality::Optional);
let bytes_multi = open_bytes_column(reader, "bytes_multi");
assert_eq!(bytes_multi.get_cardinality(), Cardinality::Multivalued);
for row_id in 0..num_docs {
let doc = row_id % STRING_BYTES_NUM_DOCS;
assert_eq!(
str_values(&str_full, row_id),
vec![format!("str-full-{doc}")]
);
assert_eq!(
str_values(&str_optional, row_id),
if doc.is_multiple_of(2) {
vec![format!("str-optional-{doc}")]
} else {
Vec::new()
}
);
assert_eq!(
str_values(&str_multi, row_id),
vec![format!("str-multi-{doc}-a"), format!("str-multi-{doc}-b")]
);
assert_eq!(
bytes_values(&bytes_full, row_id),
vec![vec![doc as u8, 0, 255]]
);
assert_eq!(
bytes_values(&bytes_optional, row_id),
if doc.is_multiple_of(2) {
vec![vec![doc as u8, 1, 254]]
} else {
Vec::new()
}
);
assert_eq!(
bytes_values(&bytes_multi, row_id),
vec![vec![doc as u8, 2, 0], vec![doc as u8, 2, 255]]
);
}
}
fn open_str_column(reader: &ColumnarReader, name: &str) -> DictionaryEncodedStrColumn {
let DynamicColumn::Str(column) = reader.read_columns(name).unwrap()[0].open().unwrap() else {
panic!("expected a string column")
};
assert_eq!(column.payload_encoding(), PayloadEncoding::Dictionary);
column.as_dictionary_encoded().unwrap().clone()
}
fn open_bytes_column(reader: &ColumnarReader, name: &str) -> DictionaryEncodedBytesColumn {
let DynamicColumn::Bytes(column) = reader.read_columns(name).unwrap()[0].open().unwrap() else {
panic!("expected a byte column")
};
assert_eq!(column.payload_encoding(), PayloadEncoding::Dictionary);
column.as_dictionary_encoded().unwrap().clone()
}
fn str_values(column: &DictionaryEncodedStrColumn, row_id: u32) -> Vec<String> {
column
.term_ords(row_id)
.map(|term_ord| {
let mut output = String::new();
assert!(column.ord_to_str(term_ord, &mut output).unwrap());
output
})
.collect()
}
fn bytes_values(column: &DictionaryEncodedBytesColumn, row_id: u32) -> Vec<Vec<u8>> {
column
.term_ords(row_id)
.map(|term_ord| {
let mut output = Vec::new();
assert!(column.ord_to_bytes(term_ord, &mut output).unwrap());
output
})
.collect()
}
fn test_format(path: &str) {
let file_content = std::fs::read(path).unwrap();
let reader = ColumnarReader::open(file_content).unwrap();
+36
View File
@@ -1,5 +1,7 @@
use serde::{Deserialize, Serialize};
use crate::InvalidData;
/// Encoding used to store string and byte column payloads.
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
@@ -10,3 +12,37 @@ pub enum PayloadEncoding {
/// Store values directly without assigning them dictionary ordinals.
Plain,
}
impl PayloadEncoding {
pub(crate) fn to_code(self) -> u8 {
match self {
PayloadEncoding::Dictionary => 0,
PayloadEncoding::Plain => 1,
}
}
pub(crate) fn try_from_code(code: u8) -> Result<Self, InvalidData> {
match code {
0 => Ok(PayloadEncoding::Dictionary),
1 => Ok(PayloadEncoding::Plain),
_ => Err(InvalidData),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_payload_encoding_codes() {
let mut valid_codes = Vec::new();
for code in u8::MIN..=u8::MAX {
if let Ok(encoding) = PayloadEncoding::try_from_code(code) {
assert_eq!(encoding.to_code(), code);
valid_codes.push(code);
}
}
assert_eq!(valid_codes, [0, 1]);
}
}
+2 -2
View File
@@ -26,7 +26,7 @@ fn test_dataframe_writer_str() {
assert_eq!(columnar.num_columns(), 1);
let cols: Vec<DynamicColumnHandle> = columnar.read_columns("my_string").unwrap();
assert_eq!(cols.len(), 1);
assert_eq!(cols[0].num_bytes(), 73);
assert_eq!(cols[0].num_bytes(), 74);
}
#[test]
@@ -40,7 +40,7 @@ fn test_dataframe_writer_bytes() {
assert_eq!(columnar.num_columns(), 1);
let cols: Vec<DynamicColumnHandle> = columnar.read_columns("my_string").unwrap();
assert_eq!(cols.len(), 1);
assert_eq!(cols[0].num_bytes(), 73);
assert_eq!(cols[0].num_bytes(), 74);
}
#[test]
+44 -8
View File
@@ -6,8 +6,8 @@ Add a user-selectable payload encoding for string and byte fast fields:
```rust
pub enum PayloadEncoding {
Plain,
Dictionary,
Plain,
}
```
@@ -50,9 +50,9 @@ and on-disk format all need it. Re-export it from `tantivy::schema` for normal T
```rust
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
pub enum PayloadEncoding {
Plain,
#[default]
Dictionary,
Plain,
}
```
@@ -246,12 +246,47 @@ V3 string and byte payloads start with a stable encoding discriminant:
encoding tag | encoding-specific payload
```
For `Dictionary`, the bytes after the tag use the current dictionary payload layout unchanged.
For `Plain`, define a versioned layout containing the column index, OnPair16 model, compressed
payload, offsets, and explicit region lengths in a footer.
The one-byte encoding tags are fixed independently of the Rust enum declaration order:
The exact ordering should allow `OwnedBytes::split` operations without copying. All serialized
lengths must be checked before slicing.
```text
0 = Dictionary
1 = Plain
2..=255 = reserved; readers reject them
```
For `Dictionary`, the bytes after the tag use the V2 dictionary payload layout unchanged:
```text
0u8 | dictionary | column index and term ordinals | dictionary_num_bytes:u32 LE
```
`dictionary_num_bytes` counts only `dictionary`, excluding the encoding tag. For `Plain`, the V3
layout is:
```text
1u8
| column_index
| onpair16_model
| compressed_values
| value_offsets
| column_index_num_bytes:u32 LE
| model_num_bytes:u32 LE
| compressed_values_num_bytes:u64 LE
| value_offsets_num_bytes:u32 LE
| num_values:u32 LE
```
`value_offsets` is a serialized monotonic `u64` column with exactly `num_values + 1` entries. Its
first entry is zero, its last entry equals `compressed_values_num_bytes`, and every adjacent pair
delimits one independently compressed value. `onpair16_model` uses the codec's canonical model
serialization, which must be frozen alongside the standalone plain-column implementation.
The fixed 24-byte footer is read from the end first. The four regions are then split from left to
right without copying. Checked conversion to `usize`, checked length sums, and region bounds are
required before any `OwnedBytes::split` call; unrecognized tags and trailing bytes are invalid.
V1 and V2 never consume a tag. V3 always consumes exactly one tag byte for string and byte
payloads, including dictionary payloads.
Other column types do not need an encoding tag and retain their existing V3 representation unless
the version implementation requires a uniform envelope.
@@ -259,7 +294,8 @@ the version implementation requires a uniform envelope.
### Compatibility tests
Expand `columnar/src/compat_tests.rs` fixtures so V1 and V2 include dictionary-encoded string and
byte columns for all supported cardinalities. Tests must:
byte columns for all supported cardinalities. The historical fixtures are
`v1_string_bytes.columnar` and `v2_string_bytes.columnar`. Tests must:
- Open and read old values through the new enum variants.
- Assert that old columns become `DictionaryEncoded`.