Skip to content

Commit ee17bfd

Browse files
authored
feat(logging): bridge native IVF-PQ diagnostics to SLF4J (#75)
1 parent cf0506d commit ee17bfd

11 files changed

Lines changed: 925 additions & 15 deletions

File tree

.github/workflows/ci.yml

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -191,12 +191,14 @@ jobs:
191191
java-version: '8'
192192
distribution: 'temurin'
193193

194-
- name: Test Java API
195-
run: mvn -f java/pom.xml test
196-
197194
- name: Build JNI library
198195
run: cargo build -p paimon-vindex-jni --release
199196

197+
- name: Test Java API
198+
run: >
199+
mvn -f java/pom.xml test
200+
-Dpaimon.vindex.native.path="${{ github.workspace }}/target/release/libpaimon_vindex_jni.so"
201+
200202
- name: Test JNI native behavior
201203
run: |
202204
java -cp java/target/test-classes:java/target/classes org.apache.paimon.index.vector.VectorIndexNativeValidationTest \

core/src/io.rs

Lines changed: 66 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ use crate::pq::ProductQuantizer;
2727
use rayon::prelude::*;
2828
use std::io;
2929
use std::mem::size_of;
30+
use std::time::{Duration, Instant};
3031

3132
pub const MAGIC: u32 = 0x49565051; // "IVPQ"
3233
pub const VERSION: u32 = 1;
@@ -95,6 +96,56 @@ pub trait SeekRead: Send {
9596
}
9697
}
9798

99+
#[derive(Clone, Copy, Debug, Default)]
100+
pub(crate) struct ReadMetrics {
101+
pub elapsed: Duration,
102+
pub calls: usize,
103+
pub requested_bytes: usize,
104+
}
105+
106+
struct MeasuredSeekRead<R> {
107+
inner: R,
108+
metrics: Option<ReadMetrics>,
109+
}
110+
111+
impl<R> MeasuredSeekRead<R> {
112+
fn new(inner: R) -> Self {
113+
Self {
114+
inner,
115+
metrics: None,
116+
}
117+
}
118+
}
119+
120+
impl<R: SeekRead> SeekRead for MeasuredSeekRead<R> {
121+
fn pread(&mut self, ranges: &mut [ReadRequest<'_>]) -> io::Result<()> {
122+
let measurement = self.metrics.as_ref().map(|_| {
123+
let requested_bytes = ranges.iter().fold(0usize, |total, request| {
124+
total.saturating_add(request.buf.len())
125+
});
126+
(Instant::now(), requested_bytes)
127+
});
128+
let result = self.inner.pread(ranges);
129+
if let (Some(metrics), Some((started, requested_bytes))) =
130+
(self.metrics.as_mut(), measurement)
131+
{
132+
metrics.elapsed += started.elapsed();
133+
metrics.calls = metrics.calls.saturating_add(1);
134+
metrics.requested_bytes = metrics.requested_bytes.saturating_add(requested_bytes);
135+
}
136+
result
137+
}
138+
139+
fn try_clone_reader(&self) -> io::Result<Option<Self>> {
140+
// Clones start with metrics disabled, so their I/O is not included here.
141+
Ok(self.inner.try_clone_reader()?.map(Self::new))
142+
}
143+
144+
fn read_capabilities(&self) -> SeekReadCapabilities {
145+
self.inner.read_capabilities()
146+
}
147+
}
148+
98149
pub(crate) struct PreadCursor<'a, R: SeekRead + ?Sized> {
99150
reader: &'a mut R,
100151
pos: u64,
@@ -456,7 +507,7 @@ fn u64_to_i64(value: u64, field: &str) -> io::Result<i64> {
456507
// --- Reader ---
457508

458509
pub struct IVFPQIndexReader<R: SeekRead> {
459-
reader: R,
510+
reader: MeasuredSeekRead<R>,
460511
pub d: usize,
461512
pub nlist: usize,
462513
pub m: usize,
@@ -578,7 +629,7 @@ impl<R: SeekRead> IVFPQIndexReader<R> {
578629
let has_opq = flags & FLAG_HAS_OPQ != 0;
579630

580631
Ok(IVFPQIndexReader {
581-
reader,
632+
reader: MeasuredSeekRead::new(reader),
582633
d,
583634
nlist,
584635
m,
@@ -609,6 +660,14 @@ impl<R: SeekRead> IVFPQIndexReader<R> {
609660
})
610661
}
611662

663+
pub(crate) fn begin_read_metrics(&mut self) {
664+
self.reader.metrics = Some(ReadMetrics::default());
665+
}
666+
667+
pub(crate) fn end_read_metrics(&mut self) -> ReadMetrics {
668+
self.reader.metrics.take().unwrap_or_default()
669+
}
670+
612671
/// Load centroids, codebooks, and offset table. Called automatically on first search.
613672
pub fn ensure_loaded(&mut self) -> io::Result<()> {
614673
if self.loaded {
@@ -1396,9 +1455,12 @@ mod tests {
13961455
reader.list_id_bytes_lens[non_empty_list] > 0,
13971456
"v1 files must store id_bytes_len in the offset table"
13981457
);
1458+
let expected_requested_bytes = reader.list_payload_len(non_empty_list).unwrap();
1459+
reader.begin_read_metrics();
13991460
let mut lists = reader
14001461
.read_inverted_list_payloads(&[non_empty_list])
14011462
.unwrap();
1463+
let read_metrics = reader.end_read_metrics();
14021464
let list = lists.pop().unwrap();
14031465
let read_ids = &list.ids;
14041466
let codes = list.codes();
@@ -1416,6 +1478,8 @@ mod tests {
14161478
stats.pread_calls, 1,
14171479
"delta-varint lists with offset-table id length should use one pread"
14181480
);
1481+
assert_eq!(read_metrics.calls, 1);
1482+
assert_eq!(read_metrics.requested_bytes, expected_requested_bytes);
14191483
}
14201484

14211485
#[test]

0 commit comments

Comments
 (0)