Skip to content

Commit 22a13b6

Browse files
authored
fix(table): align pk-vector source meta with Java (#531)
1 parent 0d927e1 commit 22a13b6

7 files changed

Lines changed: 203 additions & 45 deletions

File tree

crates/paimon/src/spec/pk_vector_source.rs

Lines changed: 58 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -76,15 +76,28 @@ impl PkVectorSourceFile {
7676
/// Mirrors Java `org.apache.paimon.index.pkvector.PkVectorSourceMeta`.
7777
#[derive(Debug, Clone, PartialEq, Eq)]
7878
pub struct PkVectorSourceMeta {
79+
data_level: i32,
7980
source_files: Vec<PkVectorSourceFile>,
8081
}
8182

8283
impl PkVectorSourceMeta {
83-
pub fn new(source_files: Vec<PkVectorSourceFile>) -> crate::Result<Self> {
84+
pub fn new(data_level: i32, source_files: Vec<PkVectorSourceFile>) -> crate::Result<Self> {
85+
if data_level <= 0 {
86+
return Err(data_invalid(format!(
87+
"source meta data level must be positive: {data_level}"
88+
)));
89+
}
8490
if source_files.is_empty() {
8591
return Err(data_invalid("a vector index must reference source files"));
8692
}
87-
Ok(Self { source_files })
93+
Ok(Self {
94+
data_level,
95+
source_files,
96+
})
97+
}
98+
99+
pub fn data_level(&self) -> i32 {
100+
self.data_level
88101
}
89102

90103
pub fn source_files(&self) -> &[PkVectorSourceFile] {
@@ -136,6 +149,12 @@ impl PkVectorSourceMeta {
136149
"unsupported vector source version: {version}"
137150
)));
138151
}
152+
let data_level = cursor.read_i32_be()?;
153+
if data_level <= 0 {
154+
return Err(data_invalid(format!(
155+
"source meta data level must be positive: {data_level}"
156+
)));
157+
}
139158
let count = cursor.read_i32_be()?;
140159
if count <= 0 {
141160
return Err(data_invalid("a vector index must reference source files"));
@@ -152,7 +171,7 @@ impl PkVectorSourceMeta {
152171
"unexpected trailing bytes in vector source metadata",
153172
));
154173
}
155-
Self::new(source_files)
174+
Self::new(data_level, source_files)
156175
}
157176
}
158177

@@ -270,9 +289,10 @@ mod tests {
270289
out
271290
}
272291

273-
fn frame(files: &[(&str, i64)]) -> Vec<u8> {
292+
fn frame(data_level: i32, files: &[(&str, i64)]) -> Vec<u8> {
274293
let mut out = Vec::new();
275294
out.extend_from_slice(&1i32.to_be_bytes()); // version
295+
out.extend_from_slice(&data_level.to_be_bytes());
276296
out.extend_from_slice(&(files.len() as i32).to_be_bytes());
277297
for (name, rows) in files {
278298
out.extend_from_slice(&java_write_utf(name));
@@ -319,85 +339,99 @@ mod tests {
319339

320340
#[test]
321341
fn deserialize_single_source_file() {
322-
let bytes = frame(&[("data-abc.parquet", 100)]);
342+
let bytes = frame(1, &[("data-abc.parquet", 100)]);
323343
let meta = PkVectorSourceMeta::deserialize(&bytes).unwrap();
344+
assert_eq!(meta.data_level(), 1);
324345
assert_eq!(meta.source_files().len(), 1);
325346
assert_eq!(meta.source_files()[0].file_name(), "data-abc.parquet");
326347
assert_eq!(meta.source_files()[0].row_count(), 100);
327348
}
328349

329350
#[test]
330351
fn deserialize_multi_source_files() {
331-
let bytes = frame(&[("f0", 3), ("f1", 5)]);
352+
let bytes = frame(2, &[("f0", 3), ("f1", 5)]);
332353
let meta = PkVectorSourceMeta::deserialize(&bytes).unwrap();
354+
assert_eq!(meta.data_level(), 2);
333355
assert_eq!(meta.source_files().len(), 2);
334356
assert_eq!(meta.source_files()[1].row_count(), 5);
335357
}
336358

337359
#[test]
338360
fn deserialize_rejects_bad_version() {
339-
let mut bytes = frame(&[("f0", 1)]);
361+
let mut bytes = frame(1, &[("f0", 1)]);
340362
bytes[0..4].copy_from_slice(&2i32.to_be_bytes());
341363
assert!(PkVectorSourceMeta::deserialize(&bytes).is_err());
342364
}
343365

366+
#[test]
367+
fn deserialize_rejects_zero_or_negative_data_level() {
368+
assert!(PkVectorSourceMeta::deserialize(&frame(0, &[("f0", 1)])).is_err());
369+
assert!(PkVectorSourceMeta::deserialize(&frame(-1, &[("f0", 1)])).is_err());
370+
}
371+
344372
#[test]
345373
fn deserialize_rejects_zero_count() {
346374
let mut out = Vec::new();
347375
out.extend_from_slice(&1i32.to_be_bytes());
376+
out.extend_from_slice(&1i32.to_be_bytes());
348377
out.extend_from_slice(&0i32.to_be_bytes());
349378
assert!(PkVectorSourceMeta::deserialize(&out).is_err());
350379
}
351380

352381
#[test]
353382
fn deserialize_rejects_trailing_bytes() {
354-
let mut bytes = frame(&[("f0", 1)]);
383+
let mut bytes = frame(1, &[("f0", 1)]);
355384
bytes.push(0xFF);
356385
assert!(PkVectorSourceMeta::deserialize(&bytes).is_err());
357386
}
358387

359388
#[test]
360389
fn deserialize_rejects_truncated_input() {
361-
let bytes = frame(&[("f0", 1)]);
390+
let bytes = frame(1, &[("f0", 1)]);
362391
assert!(PkVectorSourceMeta::deserialize(&bytes[..bytes.len() - 2]).is_err());
363392
}
364393

365394
#[test]
366395
fn deserialize_rejects_negative_row_count() {
367-
let bytes = frame(&[("f0", -1)]);
396+
let bytes = frame(1, &[("f0", -1)]);
368397
assert!(PkVectorSourceMeta::deserialize(&bytes).is_err());
369398
}
370399

371400
#[test]
372401
fn new_rejects_empty() {
373-
assert!(PkVectorSourceMeta::new(Vec::new()).is_err());
402+
assert!(PkVectorSourceMeta::new(1, Vec::new()).is_err());
403+
assert!(PkVectorSourceMeta::new(
404+
0,
405+
vec![PkVectorSourceFile::new("f0".to_string(), 1).unwrap()]
406+
)
407+
.is_err());
374408
}
375409

376410
#[test]
377411
fn resolve_single_file() {
378-
let meta = PkVectorSourceMeta::deserialize(&frame(&[("f0", 3)])).unwrap();
412+
let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap();
379413
assert_eq!(meta.resolve(0).unwrap(), ("f0".to_string(), 0));
380414
assert_eq!(meta.resolve(2).unwrap(), ("f0".to_string(), 2));
381415
}
382416

383417
#[test]
384418
fn resolve_multi_file_prefix_sum_boundaries() {
385419
// f0 owns ordinals 0..=2, f1 owns 3..=7.
386-
let meta = PkVectorSourceMeta::deserialize(&frame(&[("f0", 3), ("f1", 5)])).unwrap();
420+
let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3), ("f1", 5)])).unwrap();
387421
assert_eq!(meta.resolve(2).unwrap(), ("f0".to_string(), 2)); // last of f0
388422
assert_eq!(meta.resolve(3).unwrap(), ("f1".to_string(), 0)); // first of f1
389423
assert_eq!(meta.resolve(7).unwrap(), ("f1".to_string(), 4)); // last of f1
390424
}
391425

392426
#[test]
393427
fn resolve_rejects_negative_ordinal() {
394-
let meta = PkVectorSourceMeta::deserialize(&frame(&[("f0", 3)])).unwrap();
428+
let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap();
395429
assert!(meta.resolve(-1).is_err());
396430
}
397431

398432
#[test]
399433
fn resolve_rejects_ordinal_at_or_past_total() {
400-
let meta = PkVectorSourceMeta::deserialize(&frame(&[("f0", 3)])).unwrap();
434+
let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap();
401435
assert!(meta.resolve(3).is_err()); // total == 3, valid range 0..=2
402436
}
403437

@@ -422,20 +456,24 @@ mod tests {
422456
index_field_id: 0,
423457
extra_field_ids: None,
424458
index_meta: None,
425-
source_meta: Some(frame(&[("f0", 3)])),
459+
source_meta: Some(frame(1, &[("f0", 3)])),
426460
};
427461
let parsed = PkVectorSourceMeta::from_global_index_meta(&meta).unwrap();
462+
assert_eq!(parsed.data_level(), 1);
428463
assert_eq!(parsed.source_files()[0].file_name(), "f0");
429464
}
430465

431466
#[test]
432467
fn resolve_rejects_row_count_overflow() {
433468
// Two individually-valid row counts whose prefix sum overflows i64.
434469
// Resolving past the first file forces the checked_add on the second.
435-
let meta = PkVectorSourceMeta::new(vec![
436-
PkVectorSourceFile::new("f0".to_string(), i64::MAX).unwrap(),
437-
PkVectorSourceFile::new("f1".to_string(), 1).unwrap(),
438-
])
470+
let meta = PkVectorSourceMeta::new(
471+
1,
472+
vec![
473+
PkVectorSourceFile::new("f0".to_string(), i64::MAX).unwrap(),
474+
PkVectorSourceFile::new("f1".to_string(), 1).unwrap(),
475+
],
476+
)
439477
.unwrap();
440478
assert!(meta.resolve(i64::MAX).is_err());
441479
}

crates/paimon/src/table/pk_vector_orchestrator.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -764,6 +764,7 @@ mod e2e_tests {
764764
fn ann_segment(sources: &[(&str, i64)]) -> BucketAnnSegment {
765765
BucketAnnSegment::for_test(
766766
PkVectorSourceMeta::new(
767+
1,
767768
sources
768769
.iter()
769770
.map(|(n, r)| PkVectorSourceFile::new((*n).to_string(), *r).unwrap())

0 commit comments

Comments
 (0)