This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 17a993ec fix(table): align bitmap index ordering with Java (#601)
17a993ec is described below
commit 17a993ecd9ccfa371c453fd1a5457b9ef87b537b
Author: QuakeWang <[email protected]>
AuthorDate: Thu Jul 23 20:54:40 2026 +0800
fix(table): align bitmap index ordering with Java (#601)
Bitmap lookup mixed serialized and logical key order, while floating-point
keys did not follow Java's canonical NaN contract. This could select the wrong
dictionary block and made equality pruning unnecessarily conservative.
Use a Bitmap-specific Java key codec and logical ordering, keep
residual-sensitive floating predicates fail-open, and cover the behavior with
Java and current-format fixtures.
---
.../paimon/src/table/bitmap_global_index_reader.rs | 589 +++++++++++++++++++--
.../src/table/btree_global_index_build_builder.rs | 326 +++++++++++-
crates/paimon/src/table/global_index_scanner.rs | 570 +++++++++++++++++++-
.../testdata/bitmap/bitmap_int_logical_java.index | Bin 0 -> 274 bytes
.../bitmap/bitmap_int_logical_java.index.meta | Bin 0 -> 19 bytes
5 files changed, 1400 insertions(+), 85 deletions(-)
diff --git a/crates/paimon/src/table/bitmap_global_index_reader.rs
b/crates/paimon/src/table/bitmap_global_index_reader.rs
index 0c6172b7..151aac8e 100644
--- a/crates/paimon/src/table/bitmap_global_index_reader.rs
+++ b/crates/paimon/src/table/bitmap_global_index_reader.rs
@@ -19,6 +19,7 @@
//!
//! Reference: `org.apache.paimon.globalindex.bitmap.BitmapGlobalIndexFormat`.
+use crate::btree::key_serde::KeyComparator;
use crate::btree::var_len::{decode_var_int, decode_var_long, encode_var_int,
encode_var_long};
use crate::btree::BTreeIndexMeta;
use crate::btree::{make_key_comparator, serialize_datum, BlockCompressionType};
@@ -35,6 +36,80 @@ const MAGIC: i32 = 0x4247_4958;
const VERSION: i32 = 1;
const FOOTER_LENGTH: usize = 48;
const BLOCK_TRAILER_LENGTH: usize = 5;
+const JAVA_CANONICAL_FLOAT_NAN_BITS: u32 = 0x7fc0_0000;
+const JAVA_CANONICAL_DOUBLE_NAN_BITS: u64 = 0x7ff8_0000_0000_0000;
+
+// Bitmap follows current Java's floating-point key contract. Shared BTree key
+// serde intentionally keeps the already-persisted Rust contract.
+pub(crate) fn make_bitmap_key_comparator(data_type: &DataType) ->
KeyComparator {
+ match data_type {
+ DataType::Float(_) => Box::new(|left, right| {
+ let left = f32::from_le_bytes(left[..4].try_into().unwrap());
+ let right = f32::from_le_bytes(right[..4].try_into().unwrap());
+ compare_float_like_java(left, right)
+ }),
+ DataType::Double(_) => Box::new(|left, right| {
+ let left = f64::from_le_bytes(left[..8].try_into().unwrap());
+ let right = f64::from_le_bytes(right[..8].try_into().unwrap());
+ compare_double_like_java(left, right)
+ }),
+ _ => make_key_comparator(data_type),
+ }
+}
+
+pub(crate) fn serialize_bitmap_datum(datum: &Datum, data_type: &DataType) ->
Vec<u8> {
+ match (datum, data_type) {
+ (Datum::Float(value), DataType::Float(_)) => {
+ let bits = if value.is_nan() {
+ JAVA_CANONICAL_FLOAT_NAN_BITS
+ } else {
+ value.to_bits()
+ };
+ bits.to_le_bytes().to_vec()
+ }
+ (Datum::Double(value), DataType::Double(_)) => {
+ let bits = if value.is_nan() {
+ JAVA_CANONICAL_DOUBLE_NAN_BITS
+ } else {
+ value.to_bits()
+ };
+ bits.to_le_bytes().to_vec()
+ }
+ _ => serialize_datum(datum, data_type),
+ }
+}
+
+pub(crate) fn is_bitmap_floating_residual_sensitive_op(op: PredicateOperator)
-> bool {
+ matches!(
+ op,
+ PredicateOperator::NotEq
+ | PredicateOperator::NotIn
+ | PredicateOperator::Lt
+ | PredicateOperator::LtEq
+ | PredicateOperator::Gt
+ | PredicateOperator::GtEq
+ | PredicateOperator::Between
+ | PredicateOperator::NotBetween
+ )
+}
+
+fn compare_float_like_java(left: f32, right: f32) -> Ordering {
+ match (left.is_nan(), right.is_nan()) {
+ (true, true) => Ordering::Equal,
+ (true, false) => Ordering::Greater,
+ (false, true) => Ordering::Less,
+ (false, false) => left.total_cmp(&right),
+ }
+}
+
+fn compare_double_like_java(left: f64, right: f64) -> Ordering {
+ match (left.is_nan(), right.is_nan()) {
+ (true, true) => Ordering::Equal,
+ (true, false) => Ordering::Greater,
+ (false, true) => Ordering::Less,
+ (false, false) => left.total_cmp(&right),
+ }
+}
#[derive(Clone, Copy)]
struct BlockInfo {
@@ -124,12 +199,17 @@ impl<F: Fn(&[u8], &[u8]) -> Ordering>
BitmapGlobalIndexWriter<F> {
}
pub(crate) async fn finish(mut self) -> io::Result<BitmapWriteResult> {
+ let mut bitmaps = std::mem::take(&mut self.bitmaps)
+ .into_iter()
+ .collect::<Vec<_>>();
+ bitmaps.sort_by(|(left, _), (right, _)| (self.key_comparator)(left,
right));
+
let mut bytes = Vec::new();
write_bitmap_index_bytes(
&mut bytes,
&self.null_rows,
&self.non_null_rows,
- &self.bitmaps,
+ &bitmaps,
self.dictionary_block_size,
self.compression_type,
)?;
@@ -190,58 +270,61 @@ impl BitmapGlobalIndexReader {
literals: &[Datum],
data_type: &DataType,
) -> io::Result<RoaringTreemap> {
+ if is_floating_point(data_type) &&
is_bitmap_floating_residual_sensitive_op(op) {
+ return self.is_not_null().await;
+ }
match op {
PredicateOperator::Eq => {
- let key = serialize_datum(&literals[0], data_type);
- self.equal(&key).await
+ let key = serialize_bitmap_datum(&literals[0], data_type);
+ self.equal(&key, data_type).await
}
PredicateOperator::NotEq => {
let mut result = self.is_not_null().await?;
- let key = serialize_datum(&literals[0], data_type);
- result -= self.equal(&key).await?;
+ let key = serialize_bitmap_datum(&literals[0], data_type);
+ result -= self.equal(&key, data_type).await?;
Ok(result)
}
PredicateOperator::In => {
let keys = literals
.iter()
- .map(|literal| serialize_datum(literal, data_type))
+ .map(|literal| serialize_bitmap_datum(literal, data_type))
.collect::<Vec<_>>();
- self.in_keys(&keys).await
+ self.in_keys(&keys, data_type).await
}
PredicateOperator::NotIn => {
let mut result = self.is_not_null().await?;
let keys = literals
.iter()
- .map(|literal| serialize_datum(literal, data_type))
+ .map(|literal| serialize_bitmap_datum(literal, data_type))
.collect::<Vec<_>>();
- result -= self.in_keys(&keys).await?;
+ result -= self.in_keys(&keys, data_type).await?;
Ok(result)
}
PredicateOperator::IsNull => self.is_null().await,
PredicateOperator::IsNotNull => self.is_not_null().await,
PredicateOperator::Lt => {
- let key = serialize_datum(&literals[0], data_type);
+ let key = serialize_bitmap_datum(&literals[0], data_type);
self.scan_dictionary(data_type, |candidate, cmp|
cmp(candidate, &key).is_lt())
.await
}
PredicateOperator::LtEq => {
- let key = serialize_datum(&literals[0], data_type);
+ let key = serialize_bitmap_datum(&literals[0], data_type);
self.scan_dictionary(data_type, |candidate, cmp|
!cmp(candidate, &key).is_gt())
.await
}
PredicateOperator::Gt => {
- let key = serialize_datum(&literals[0], data_type);
+ let key = serialize_bitmap_datum(&literals[0], data_type);
self.scan_dictionary(data_type, |candidate, cmp|
cmp(candidate, &key).is_gt())
.await
}
PredicateOperator::GtEq => {
- let key = serialize_datum(&literals[0], data_type);
+ let key = serialize_bitmap_datum(&literals[0], data_type);
self.scan_dictionary(data_type, |candidate, cmp|
!cmp(candidate, &key).is_lt())
.await
}
PredicateOperator::Between => {
- let from = serialize_datum(&literals[0], data_type);
- let to = serialize_datum(&literals[1], data_type);
+ let from = serialize_bitmap_datum(&literals[0], data_type);
+ let to = serialize_bitmap_datum(&literals[1], data_type);
self.scan_dictionary(data_type, |candidate, cmp| {
!cmp(candidate, &from).is_lt() && !cmp(candidate,
&to).is_gt()
})
@@ -249,8 +332,8 @@ impl BitmapGlobalIndexReader {
}
PredicateOperator::NotBetween => {
let mut result = self.is_not_null().await?;
- let from = serialize_datum(&literals[0], data_type);
- let to = serialize_datum(&literals[1], data_type);
+ let from = serialize_bitmap_datum(&literals[0], data_type);
+ let to = serialize_bitmap_datum(&literals[1], data_type);
let inside = self
.scan_dictionary(data_type, |candidate, cmp| {
!cmp(candidate, &from).is_lt() && !cmp(candidate,
&to).is_gt()
@@ -266,7 +349,7 @@ impl BitmapGlobalIndexReader {
"Bitmap global index starts_with only supports string
columns",
));
}
- let prefix = serialize_datum(&literals[0], data_type);
+ let prefix = serialize_bitmap_datum(&literals[0], data_type);
if prefix.is_empty() {
return self.is_not_null().await;
}
@@ -280,7 +363,7 @@ impl BitmapGlobalIndexReader {
"Bitmap global index ends_with only supports string
columns",
));
}
- let suffix = serialize_datum(&literals[0], data_type);
+ let suffix = serialize_bitmap_datum(&literals[0], data_type);
if suffix.is_empty() {
return self.is_not_null().await;
}
@@ -294,7 +377,7 @@ impl BitmapGlobalIndexReader {
"Bitmap global index contains only supports string
columns",
));
}
- let needle = serialize_datum(&literals[0], data_type);
+ let needle = serialize_bitmap_datum(&literals[0], data_type);
if needle.is_empty() {
return self.is_not_null().await;
}
@@ -325,6 +408,9 @@ impl BitmapGlobalIndexReader {
from_inclusive: bool,
to_inclusive: bool,
) -> io::Result<RoaringTreemap> {
+ if is_floating_point(data_type) {
+ return self.is_not_null().await;
+ }
self.scan_dictionary(data_type, |candidate, cmp| {
let from_cmp = cmp(candidate, from);
let to_cmp = cmp(candidate, to);
@@ -342,23 +428,33 @@ impl BitmapGlobalIndexReader {
self.read_bitmap(self.footer.non_null_rows_block).await
}
- async fn equal(&self, key: &[u8]) -> io::Result<RoaringTreemap> {
- match self.find_bitmap_block(key).await? {
+ async fn equal(&self, key: &[u8], data_type: &DataType) ->
io::Result<RoaringTreemap> {
+ let logical_cmp = make_bitmap_key_comparator(data_type);
+ self.equal_with_comparator(key, logical_cmp.as_ref()).await
+ }
+
+ async fn equal_with_comparator(
+ &self,
+ key: &[u8],
+ logical_cmp: &(dyn Fn(&[u8], &[u8]) -> Ordering + Send + Sync),
+ ) -> io::Result<RoaringTreemap> {
+ match self.find_bitmap_block(key, logical_cmp).await? {
Some(block) => self.read_bitmap(block).await,
None => Ok(RoaringTreemap::new()),
}
}
- async fn in_keys(&self, keys: &[Vec<u8>]) -> io::Result<RoaringTreemap> {
+ async fn in_keys(&self, keys: &[Vec<u8>], data_type: &DataType) ->
io::Result<RoaringTreemap> {
let mut sorted_keys = keys.to_vec();
sorted_keys.sort();
sorted_keys.dedup();
+ let logical_cmp = make_bitmap_key_comparator(data_type);
let mut result = RoaringTreemap::new();
for key in sorted_keys {
- if let Some(block) = self.find_bitmap_block(&key).await? {
- result |= self.read_bitmap(block).await?;
- }
+ result |= self
+ .equal_with_comparator(&key, logical_cmp.as_ref())
+ .await?;
}
Ok(result)
}
@@ -368,7 +464,7 @@ impl BitmapGlobalIndexReader {
data_type: &DataType,
predicate: impl Fn(&[u8], &dyn Fn(&[u8], &[u8]) -> Ordering) -> bool,
) -> io::Result<RoaringTreemap> {
- let cmp = make_key_comparator(data_type);
+ let cmp = make_bitmap_key_comparator(data_type);
self.scan_serialized_dictionary(|candidate| predicate(candidate,
cmp.as_ref()))
.await
}
@@ -388,13 +484,16 @@ impl BitmapGlobalIndexReader {
Ok(result)
}
- async fn find_bitmap_block(&self, key: &[u8]) ->
io::Result<Option<BlockInfo>> {
- let Some(block_meta) = self.find_dictionary_block_meta(key) else {
+ async fn find_bitmap_block(
+ &self,
+ key: &[u8],
+ logical_cmp: &(dyn Fn(&[u8], &[u8]) -> Ordering + Send + Sync),
+ ) -> io::Result<Option<BlockInfo>> {
+ let Some(block_meta) = self.find_dictionary_block_meta(key,
logical_cmp) else {
return Ok(None);
};
-
for entry in self.read_dictionary_block(block_meta.block).await? {
- match compare_unsigned(&entry.key, key) {
+ match logical_cmp(&entry.key, key) {
Ordering::Equal => return Ok(Some(entry.bitmap_block)),
Ordering::Greater => return Ok(None),
Ordering::Less => {}
@@ -403,7 +502,11 @@ impl BitmapGlobalIndexReader {
Ok(None)
}
- fn find_dictionary_block_meta(&self, key: &[u8]) ->
Option<&DictionaryBlockMeta> {
+ fn find_dictionary_block_meta(
+ &self,
+ key: &[u8],
+ compare: impl Fn(&[u8], &[u8]) -> Ordering,
+ ) -> Option<&DictionaryBlockMeta> {
if self.dictionary_blocks.is_empty() {
return None;
}
@@ -411,7 +514,7 @@ impl BitmapGlobalIndexReader {
let mut high = self.dictionary_blocks.len();
while low < high {
let mid = (low + high) / 2;
- if compare_unsigned(&self.dictionary_blocks[mid].first_key, key)
!= Ordering::Greater {
+ if compare(&self.dictionary_blocks[mid].first_key, key) !=
Ordering::Greater {
low = mid + 1;
} else {
high = mid;
@@ -519,7 +622,6 @@ fn read_index_block(bytes: &[u8]) ->
io::Result<Vec<DictionaryBlockMeta>> {
block: block_info(offset, length)?,
});
}
- blocks.sort_by(|left, right| compare_unsigned(&left.first_key,
&right.first_key));
Ok(blocks)
}
@@ -627,20 +729,14 @@ fn read_i32_be(bytes: &[u8], offset: usize) ->
io::Result<i32> {
Ok(i32::from_be_bytes(value.try_into().unwrap()))
}
-fn compare_unsigned(left: &[u8], right: &[u8]) -> Ordering {
- for (left, right) in left.iter().zip(right.iter()) {
- match left.cmp(right) {
- Ordering::Equal => {}
- non_eq => return non_eq,
- }
- }
- left.len().cmp(&right.len())
-}
-
fn is_character_string(data_type: &DataType) -> bool {
matches!(data_type, DataType::Char(_) | DataType::VarChar(_))
}
+fn is_floating_point(data_type: &DataType) -> bool {
+ matches!(data_type, DataType::Float(_) | DataType::Double(_))
+}
+
fn string_literal(literals: &[Datum], op: PredicateOperator) ->
io::Result<&str> {
match literals.first() {
Some(Datum::String(value)) => Ok(value),
@@ -666,7 +762,7 @@ fn write_bitmap_index_bytes(
out: &mut Vec<u8>,
null_rows: &RoaringTreemap,
non_null_rows: &RoaringTreemap,
- bitmaps: &BTreeMap<Vec<u8>, RoaringTreemap>,
+ bitmaps: &[(Vec<u8>, RoaringTreemap)],
dictionary_block_size: usize,
compression_type: BlockCompressionType,
) -> io::Result<()> {
@@ -706,7 +802,7 @@ fn write_bitmap_block(out: &mut Vec<u8>, bitmap:
&RoaringTreemap) -> io::Result<
fn write_dictionary_and_bitmap_blocks(
out: &mut Vec<u8>,
- bitmaps: &BTreeMap<Vec<u8>, RoaringTreemap>,
+ bitmaps: &[(Vec<u8>, RoaringTreemap)],
dictionary_block_size: usize,
compression_type: BlockCompressionType,
) -> io::Result<(Vec<DictionaryBlockMeta>, usize)> {
@@ -890,3 +986,404 @@ fn u64_to_i64(value: u64) -> io::Result<i64> {
)
})
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::btree::test_util::{BytesFileRead, VecFileWrite};
+ use crate::spec::{DoubleType, FloatType, IntType};
+ use std::ops::Range;
+ use std::sync::{Arc, Mutex};
+
+ struct TrackingFileRead {
+ bytes: Bytes,
+ ranges: Arc<Mutex<Vec<Range<u64>>>>,
+ }
+
+ #[async_trait::async_trait]
+ impl FileRead for TrackingFileRead {
+ async fn read(&self, range: Range<u64>) -> crate::Result<Bytes> {
+ self.ranges.lock().unwrap().push(range.clone());
+ Ok(self.bytes.slice(range.start as usize..range.end as usize))
+ }
+ }
+
+ async fn write_bitmap_bytes(
+ data_type: &DataType,
+ values: &[(Datum, i64)],
+ dictionary_block_size: usize,
+ ) -> Vec<u8> {
+ let output = VecFileWrite::new();
+ let captured = output.clone();
+ let mut writer = BitmapGlobalIndexWriter::new(
+ Box::new(output),
+ dictionary_block_size,
+ BlockCompressionType::None,
+ make_bitmap_key_comparator(data_type),
+ );
+ for (value, row_id) in values {
+ let key = serialize_bitmap_datum(value, data_type);
+ writer.write(Some(&key), *row_id).unwrap();
+ }
+ writer.finish().await.unwrap();
+ captured.to_vec()
+ }
+
+ async fn write_and_read_entries(
+ data_type: &DataType,
+ values: &[(Datum, i64)],
+ ) -> (BitmapGlobalIndexReader, Vec<DictionaryEntry>) {
+ let bytes = write_bitmap_bytes(data_type, values, 1 << 20).await;
+ let reader = BitmapGlobalIndexReader::open(
+ Box::new(BytesFileRead(Bytes::from(bytes.clone()))),
+ bytes.len() as u64,
+ )
+ .await
+ .unwrap();
+ let mut entries = Vec::new();
+ for block in &reader.dictionary_blocks {
+
entries.extend(reader.read_dictionary_block(block.block).await.unwrap());
+ }
+ (reader, entries)
+ }
+
+ #[test]
+ fn test_bitmap_floating_residual_sensitive_operator_set() {
+ for op in [
+ PredicateOperator::NotEq,
+ PredicateOperator::NotIn,
+ PredicateOperator::Lt,
+ PredicateOperator::LtEq,
+ PredicateOperator::Gt,
+ PredicateOperator::GtEq,
+ PredicateOperator::Between,
+ PredicateOperator::NotBetween,
+ ] {
+ assert!(is_bitmap_floating_residual_sensitive_op(op), "{op}");
+ }
+ for op in [
+ PredicateOperator::Eq,
+ PredicateOperator::In,
+ PredicateOperator::IsNull,
+ PredicateOperator::IsNotNull,
+ PredicateOperator::StartsWith,
+ PredicateOperator::EndsWith,
+ PredicateOperator::Contains,
+ PredicateOperator::Like,
+ ] {
+ assert!(!is_bitmap_floating_residual_sensitive_op(op), "{op}");
+ }
+ }
+
+ #[tokio::test]
+ async fn test_writer_orders_numeric_keys_logically() {
+ let data_type = DataType::Int(IntType::new());
+ let output = VecFileWrite::new();
+ let captured = output.clone();
+ let mut writer = BitmapGlobalIndexWriter::new(
+ Box::new(output),
+ 1 << 20,
+ BlockCompressionType::None,
+ make_bitmap_key_comparator(&data_type),
+ );
+
+ writer.write(None, 5).unwrap();
+ for (value, row_id) in [(-1i32, 0), (0, 1), (0, 2), (1, 3), (256, 4)] {
+ writer.write(Some(&value.to_le_bytes()), row_id).unwrap();
+ }
+ writer.finish().await.unwrap();
+
+ let bytes = captured.to_vec();
+ let reader = BitmapGlobalIndexReader::open(
+ Box::new(BytesFileRead(Bytes::from(bytes.clone()))),
+ bytes.len() as u64,
+ )
+ .await
+ .unwrap();
+ assert_eq!(reader.dictionary_blocks.len(), 1);
+ let keys = reader
+ .read_dictionary_block(reader.dictionary_blocks[0].block)
+ .await
+ .unwrap()
+ .into_iter()
+ .map(|entry| i32::from_le_bytes(entry.key.try_into().unwrap()))
+ .collect::<Vec<_>>();
+ assert_eq!(keys, vec![-1, 0, 1, 256]);
+ }
+
+ #[tokio::test]
+ async fn test_writer_orders_and_canonicalizes_float_keys_like_java() {
+ let data_type = DataType::Float(FloatType::new());
+ let values = [
+ (Datum::Float(f32::from_bits(0xffc0_0001)), 0),
+ (Datum::Float(-1.0), 1),
+ (Datum::Float(-0.0), 2),
+ (Datum::Float(0.0), 3),
+ (Datum::Float(1.0), 4),
+ (Datum::Float(f32::from_bits(0x7fc0_0010)), 5),
+ (Datum::Float(f32::from_bits(0x7fff_1234)), 6),
+ ];
+ let (reader, entries) = write_and_read_entries(&data_type,
&values).await;
+
+ let key_bits = entries
+ .iter()
+ .map(|entry|
u32::from_le_bytes(entry.key.as_slice().try_into().unwrap()))
+ .collect::<Vec<_>>();
+ assert_eq!(
+ key_bits,
+ vec![
+ (-1.0f32).to_bits(),
+ (-0.0f32).to_bits(),
+ 0.0f32.to_bits(),
+ 1.0f32.to_bits(),
+ 0x7fc0_0000,
+ ]
+ );
+ let nan_rows = reader
+ .read_bitmap(entries.last().unwrap().bitmap_block)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(nan_rows, vec![0, 5, 6]);
+ }
+
+ #[tokio::test]
+ async fn test_writer_orders_and_canonicalizes_double_keys_like_java() {
+ let data_type = DataType::Double(DoubleType::new());
+ let values = [
+ (Datum::Double(f64::from_bits(0xfff8_0000_0000_0001)), 0),
+ (Datum::Double(-1.0), 1),
+ (Datum::Double(-0.0), 2),
+ (Datum::Double(0.0), 3),
+ (Datum::Double(1.0), 4),
+ (Datum::Double(f64::from_bits(0x7ff8_0000_0000_0010)), 5),
+ (Datum::Double(f64::from_bits(0x7fff_1234_5678_9abc)), 6),
+ ];
+ let (reader, entries) = write_and_read_entries(&data_type,
&values).await;
+
+ let key_bits = entries
+ .iter()
+ .map(|entry|
u64::from_le_bytes(entry.key.as_slice().try_into().unwrap()))
+ .collect::<Vec<_>>();
+ assert_eq!(
+ key_bits,
+ vec![
+ (-1.0f64).to_bits(),
+ (-0.0f64).to_bits(),
+ 0.0f64.to_bits(),
+ 1.0f64.to_bits(),
+ 0x7ff8_0000_0000_0000,
+ ]
+ );
+ let nan_rows = reader
+ .read_bitmap(entries.last().unwrap().bitmap_block)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(nan_rows, vec![0, 5, 6]);
+ }
+
+ async fn assert_eq_and_in_read_only_logical_candidate_block(
+ data_type: DataType,
+ negative_nan: Datum,
+ positive_nan: Datum,
+ canonical_nan: Datum,
+ negative_one: Datum,
+ zero: Datum,
+ one: Datum,
+ ) {
+ let values = [
+ (negative_nan.clone(), 0),
+ (positive_nan.clone(), 1),
+ (canonical_nan.clone(), 2),
+ (negative_one.clone(), 3),
+ (zero, 4),
+ (one, 5),
+ ];
+ let bytes = write_bitmap_bytes(&data_type, &values, 1).await;
+ let ranges = Arc::new(Mutex::new(Vec::new()));
+ let reader = BitmapGlobalIndexReader::open(
+ Box::new(TrackingFileRead {
+ bytes: Bytes::from(bytes.clone()),
+ ranges: Arc::clone(&ranges),
+ }),
+ bytes.len() as u64,
+ )
+ .await
+ .unwrap();
+ assert!(reader.dictionary_blocks.len() > 1);
+
+ let cases = [
+ (PredicateOperator::Eq, vec![negative_nan], vec![0, 1, 2]),
+ (
+ PredicateOperator::In,
+ vec![positive_nan, canonical_nan],
+ vec![0, 1, 2],
+ ),
+ (PredicateOperator::Eq, vec![negative_one], vec![3]),
+ ];
+ for (op, literals, expected) in cases {
+ let key = serialize_bitmap_datum(&literals[0], &data_type);
+ let cmp = make_bitmap_key_comparator(&data_type);
+ let expected_block = reader
+ .find_dictionary_block_meta(&key, cmp.as_ref())
+ .unwrap()
+ .block;
+ ranges.lock().unwrap().clear();
+
+ let actual = reader
+ .query(op, &literals, &data_type)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(actual, expected, "{data_type:?}: {op}");
+
+ let actual_ranges = ranges.lock().unwrap().clone();
+ assert_eq!(actual_ranges.len(), 2, "{data_type:?}: {op}");
+ assert_eq!(
+ actual_ranges[0],
+ expected_block.offset
+ ..expected_block.offset + (expected_block.length +
BLOCK_TRAILER_LENGTH) as u64,
+ "{data_type:?}: {op}"
+ );
+ }
+ }
+
+ #[tokio::test]
+ async fn test_eq_and_in_read_only_logical_candidate_block() {
+ assert_eq_and_in_read_only_logical_candidate_block(
+ DataType::Float(FloatType::new()),
+ Datum::Float(f32::from_bits(0xffc0_0001)),
+ Datum::Float(f32::from_bits(0x7fc0_0010)),
+ Datum::Float(f32::NAN),
+ Datum::Float(-1.0),
+ Datum::Float(0.0),
+ Datum::Float(1.0),
+ )
+ .await;
+ assert_eq_and_in_read_only_logical_candidate_block(
+ DataType::Double(DoubleType::new()),
+ Datum::Double(f64::from_bits(0xfff8_0000_0000_0001)),
+ Datum::Double(f64::from_bits(0x7ff8_0000_0000_0010)),
+ Datum::Double(f64::NAN),
+ Datum::Double(-1.0),
+ Datum::Double(0.0),
+ Datum::Double(1.0),
+ )
+ .await;
+ }
+
+ async fn
assert_floating_residual_sensitive_ops_return_all_non_null_candidates(
+ data_type: DataType,
+ negative_nan: Datum,
+ zero: Datum,
+ one: Datum,
+ ) {
+ let output = VecFileWrite::new();
+ let captured = output.clone();
+ let mut writer = BitmapGlobalIndexWriter::new(
+ Box::new(output),
+ 1,
+ BlockCompressionType::None,
+ make_bitmap_key_comparator(&data_type),
+ );
+ for (row_id, value) in [&negative_nan, &zero,
&one].into_iter().enumerate() {
+ let key = serialize_bitmap_datum(value, &data_type);
+ writer.write(Some(&key), row_id as i64).unwrap();
+ }
+ writer.write(None, 3).unwrap();
+ writer.finish().await.unwrap();
+
+ let bytes = captured.to_vec();
+ let reader = BitmapGlobalIndexReader::open(
+ Box::new(BytesFileRead(Bytes::from(bytes.clone()))),
+ bytes.len() as u64,
+ )
+ .await
+ .unwrap();
+ let expected = vec![0, 1, 2];
+ let cases = [
+ (PredicateOperator::NotEq, vec![negative_nan.clone()]),
+ (
+ PredicateOperator::NotIn,
+ vec![negative_nan.clone(), zero.clone()],
+ ),
+ (PredicateOperator::Lt, vec![zero.clone()]),
+ (PredicateOperator::LtEq, vec![zero.clone()]),
+ (PredicateOperator::Gt, vec![negative_nan.clone()]),
+ (PredicateOperator::GtEq, vec![negative_nan.clone()]),
+ (
+ PredicateOperator::Between,
+ vec![negative_nan.clone(), zero.clone()],
+ ),
+ (
+ PredicateOperator::NotBetween,
+ vec![negative_nan.clone(), zero.clone()],
+ ),
+ ];
+ for (op, literals) in cases {
+ let actual = reader
+ .query(op, &literals, &data_type)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(actual, expected, "{data_type:?}: {op}");
+ }
+
+ let direct_cases = [
+ (PredicateOperator::Eq, vec![negative_nan.clone()], vec![0]),
+ (
+ PredicateOperator::In,
+ vec![negative_nan.clone(), zero.clone()],
+ vec![0, 1],
+ ),
+ (PredicateOperator::Eq, vec![zero.clone()], vec![1]),
+ (
+ PredicateOperator::In,
+ vec![zero.clone(), one.clone()],
+ vec![1, 2],
+ ),
+ ];
+ for (op, literals, expected) in direct_cases {
+ let actual = reader
+ .query(op, &literals, &data_type)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(actual, expected, "{data_type:?}: direct {op}");
+ }
+
+ let from = serialize_bitmap_datum(&negative_nan, &data_type);
+ let to = serialize_bitmap_datum(&zero, &data_type);
+ let actual = reader
+ .range_query(&from, &to, &data_type, true, true)
+ .await
+ .unwrap()
+ .iter()
+ .collect::<Vec<_>>();
+ assert_eq!(actual, expected, "{data_type:?}: combined range");
+ }
+
+ #[tokio::test]
+ async fn
test_floating_residual_sensitive_ops_return_all_non_null_candidates() {
+ assert_floating_residual_sensitive_ops_return_all_non_null_candidates(
+ DataType::Float(FloatType::new()),
+ Datum::Float(f32::from_bits(0xffc0_0001)),
+ Datum::Float(0.0),
+ Datum::Float(1.0),
+ )
+ .await;
+ assert_floating_residual_sensitive_ops_return_all_non_null_candidates(
+ DataType::Double(DoubleType::new()),
+ Datum::Double(f64::from_bits(0xfff8_0000_0000_0001)),
+ Datum::Double(0.0),
+ Datum::Double(1.0),
+ )
+ .await;
+ }
+}
diff --git a/crates/paimon/src/table/btree_global_index_build_builder.rs
b/crates/paimon/src/table/btree_global_index_build_builder.rs
index b97204a7..da0b6081 100644
--- a/crates/paimon/src/table/btree_global_index_build_builder.rs
+++ b/crates/paimon/src/table/btree_global_index_build_builder.rs
@@ -15,14 +15,17 @@
// specific language governing permissions and limitations
// under the License.
-use super::bitmap_global_index_reader::{BitmapGlobalIndexWriter,
BitmapWriteResult};
+use super::bitmap_global_index_reader::{
+ make_bitmap_key_comparator, serialize_bitmap_datum,
BitmapGlobalIndexWriter, BitmapWriteResult,
+};
use super::global_index_types::{
normalize_sorted_global_index_type, BITMAP_GLOBAL_INDEX_TYPE,
BTREE_GLOBAL_INDEX_TYPE,
};
+use crate::btree::key_serde::KeyComparator;
use crate::btree::{make_key_comparator, serialize_datum, BTreeIndexWriter,
BlockCompressionType};
use crate::spec::{
bucket_dir_name, extract_datum_from_arrow, BinaryRow, CoreOptions,
DataField, DataFileMeta,
- DataType, FileKind, GlobalIndexMeta, IndexFileMeta, ROW_ID_FIELD_NAME,
+ DataType, Datum, FileKind, GlobalIndexMeta, IndexFileMeta,
ROW_ID_FIELD_NAME,
};
use crate::table::source::exclude_row_ranges;
use crate::table::source::is_data_evolution_normal_file;
@@ -41,6 +44,18 @@ const BTREE_BLOCK_SIZE: usize = 4 * 1024;
const BITMAP_DICTIONARY_BLOCK_SIZE: usize = 16 * 1024;
type BTreeKeyRow = (Option<Vec<u8>>, i64);
+type SerializeKeyFn = fn(&Datum, &DataType) -> Vec<u8>;
+
+fn make_index_key_codec(index_type: &str, data_type: &DataType) ->
(KeyComparator, SerializeKeyFn) {
+ match index_type {
+ BTREE_GLOBAL_INDEX_TYPE => (make_key_comparator(data_type),
serialize_datum),
+ BITMAP_GLOBAL_INDEX_TYPE => (
+ make_bitmap_key_comparator(data_type),
+ serialize_bitmap_datum,
+ ),
+ _ => unreachable!("normalized sorted global index type"),
+ }
+}
pub struct BTreeGlobalIndexBuildBuilder<'a> {
table: &'a Table,
@@ -189,8 +204,9 @@ impl<'a> BTreeGlobalIndexBuildBuilder<'a> {
}
})?;
let row_count = checked_row_count(shard.row_range_start,
shard.row_range_end)?;
- let mut rows = extract_index_rows(self.table, shard, index_column,
index_field).await?;
- let cmp = make_key_comparator(index_field.data_type());
+ let (cmp, serialize_key) = make_index_key_codec(index_type,
index_field.data_type());
+ let mut rows =
+ extract_index_rows(self.table, shard, index_column, index_field,
serialize_key).await?;
sort_index_rows(&mut rows, &cmp);
self.table
@@ -234,7 +250,6 @@ impl<'a> BTreeGlobalIndexBuildBuilder<'a> {
(write_result.row_count, write_result.meta)
}
BITMAP_GLOBAL_INDEX_TYPE => {
- let cmp = make_key_comparator(index_field.data_type());
let mut writer = BitmapGlobalIndexWriter::new(
writer,
BITMAP_DICTIONARY_BLOCK_SIZE,
@@ -597,6 +612,7 @@ async fn extract_index_rows(
shard: &BTreeGlobalIndexShard,
index_column: &str,
index_field: &DataField,
+ serialize_key: SerializeKeyFn,
) -> Result<Vec<BTreeKeyRow>> {
let splits = build_read_splits_for_shard(shard)?;
@@ -613,6 +629,7 @@ async fn extract_index_rows(
shard.row_range_start,
shard.row_range_end,
)?),
+ serialize_key,
)
}
@@ -655,6 +672,7 @@ fn extract_index_rows_from_batches(
data_type: &DataType,
row_range_start: i64,
expected_row_count: i64,
+ serialize_key: SerializeKeyFn,
) -> Result<Vec<BTreeKeyRow>> {
let row_count = batches.iter().map(RecordBatch::num_rows).sum::<usize>();
let mut rows = Vec::with_capacity(row_count);
@@ -705,7 +723,7 @@ fn extract_index_rows_from_batches(
expected_row_id += 1;
let key = extract_datum_from_arrow(batch, row, value_index,
data_type)?
- .map(|datum| serialize_datum(&datum, data_type));
+ .map(|datum| serialize_key(&datum, data_type));
rows.push((key, row_id - row_range_start));
}
}
@@ -768,12 +786,13 @@ mod tests {
use crate::io::FileIOBuilder;
use crate::spec::stats::BinaryTableStats;
use crate::spec::{
- BinaryType, GlobalIndexSearchMode, IndexManifest, IntType,
ManifestEntry, PredicateBuilder,
- Schema, TableSchema, VarBinaryType, VarCharType,
+ BinaryType, DoubleType, FloatType, GlobalIndexSearchMode,
IndexManifest, IntType,
+ ManifestEntry, Predicate, PredicateBuilder, Schema, TableSchema,
VarBinaryType,
+ VarCharType,
};
use crate::table::global_index_scanner::{evaluate_global_index,
GlobalIndexEvaluation};
use crate::table::{merge_row_ranges, SnapshotManager, TableCommit,
TableWrite};
- use arrow_array::{ArrayRef, Int32Array, Int64Array, StringArray};
+ use arrow_array::{ArrayRef, Float32Array, Float64Array, Int32Array,
Int64Array, StringArray};
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
use chrono::{DateTime, Utc};
use std::sync::Arc;
@@ -1056,9 +1075,15 @@ mod tests {
vec![Some(10), None, Some(30)],
vec![Some(5), Some(6), Some(7)],
);
- let rows =
- extract_index_rows_from_batches(&[batch], "id",
&DataType::Int(IntType::new()), 5, 3)
- .unwrap();
+ let rows = extract_index_rows_from_batches(
+ &[batch],
+ "id",
+ &DataType::Int(IntType::new()),
+ 5,
+ 3,
+ serialize_datum,
+ )
+ .unwrap();
assert_eq!(
rows,
@@ -1070,12 +1095,57 @@ mod tests {
);
}
+ #[test]
+ fn test_index_key_codec_scopes_java_nan_semantics_to_bitmap() {
+ fn assert_codec(
+ data_type: DataType,
+ negative_nan: Datum,
+ raw_nan_key: Vec<u8>,
+ canonical_nan_key: Vec<u8>,
+ zero: Datum,
+ ) {
+ let (btree_cmp, btree_serialize) =
+ make_index_key_codec(BTREE_GLOBAL_INDEX_TYPE, &data_type);
+ let btree_nan_key = btree_serialize(&negative_nan, &data_type);
+ let zero_key = btree_serialize(&zero, &data_type);
+ assert_eq!(btree_nan_key, raw_nan_key);
+ assert!(btree_cmp(&btree_nan_key, &zero_key).is_lt());
+
+ let (bitmap_cmp, bitmap_serialize) =
+ make_index_key_codec(BITMAP_GLOBAL_INDEX_TYPE, &data_type);
+ let bitmap_nan_key = bitmap_serialize(&negative_nan, &data_type);
+ assert_eq!(bitmap_nan_key, canonical_nan_key);
+ assert!(bitmap_cmp(&bitmap_nan_key, &zero_key).is_gt());
+ }
+
+ assert_codec(
+ DataType::Float(FloatType::new()),
+ Datum::Float(f32::from_bits(0xffc0_0001)),
+ 0xffc0_0001u32.to_le_bytes().to_vec(),
+ 0x7fc0_0000u32.to_le_bytes().to_vec(),
+ Datum::Float(0.0),
+ );
+ assert_codec(
+ DataType::Double(DoubleType::new()),
+ Datum::Double(f64::from_bits(0xfff8_0000_0000_0001)),
+ 0xfff8_0000_0000_0001u64.to_le_bytes().to_vec(),
+ 0x7ff8_0000_0000_0000u64.to_le_bytes().to_vec(),
+ Datum::Double(0.0),
+ );
+ }
+
#[test]
fn test_extract_index_rows_rejects_row_id_gap() {
let batch = index_batch(vec![Some(10), Some(30)], vec![Some(5),
Some(7)]);
- let err =
- extract_index_rows_from_batches(&[batch], "id",
&DataType::Int(IntType::new()), 5, 2)
- .expect_err("row-id gap should fail");
+ let err = extract_index_rows_from_batches(
+ &[batch],
+ "id",
+ &DataType::Int(IntType::new()),
+ 5,
+ 2,
+ serialize_datum,
+ )
+ .expect_err("row-id gap should fail");
assert!(
matches!(err, Error::DataInvalid { message, .. } if
message.contains("expected _ROW_ID"))
@@ -1126,6 +1196,7 @@ mod tests {
&DataType::VarChar(VarCharType::string_type()),
10,
2,
+ serialize_datum,
)
.unwrap();
@@ -1160,6 +1231,34 @@ mod tests {
.unwrap();
}
+ async fn scan_ids(table: &Table, predicate: Predicate) -> Vec<i32> {
+ let mut builder = table.new_read_builder();
+ builder.with_filter(predicate);
+ let plan = builder.new_scan().plan().await.unwrap();
+ let read = builder.new_read().unwrap();
+ let batches = read
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ let mut ids = batches
+ .iter()
+ .flat_map(|batch| {
+ batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap()
+ .values()
+ .iter()
+ .copied()
+ })
+ .collect::<Vec<_>>();
+ ids.sort_unstable();
+ ids
+ }
+
#[tokio::test]
async fn test_execute_writes_btree_index_manifest_and_file() {
let table_path = "memory:/test_btree_global_index_builder_e2e";
@@ -1372,6 +1471,203 @@ mod tests {
assert_eq!(row_ranges, vec![RowRange::new(0, 0), RowRange::new(2, 2)]);
}
+ #[tokio::test]
+ async fn test_bitmap_floating_candidates_preserve_residual_results() {
+ let table_path = "memory:/test_bitmap_floating_residual_candidates";
+ let schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("f", DataType::Float(FloatType::new()))
+ .column("d", DataType::Double(DoubleType::new()))
+ .options(table_options("100"))
+ .build()
+ .unwrap();
+ let table = Table::new(
+ FileIOBuilder::new("memory").build().unwrap(),
+ Identifier::new("default",
"test_bitmap_floating_residual_candidates"),
+ table_path.to_string(),
+ TableSchema::new(0, &schema),
+ None,
+ );
+ setup_dirs(&table).await;
+
+ let float_negative_nan = f32::from_bits(0xffc0_0001);
+ let double_negative_nan = f64::from_bits(0xfff8_0000_0000_0001);
+ let arrow_schema = Arc::new(ArrowSchema::new(vec![
+ ArrowField::new("id", ArrowDataType::Int32, false),
+ ArrowField::new("f", ArrowDataType::Float32, true),
+ ArrowField::new("d", ArrowDataType::Float64, true),
+ ]));
+ let batch = RecordBatch::try_new(
+ arrow_schema,
+ vec![
+ Arc::new(Int32Array::from_iter_values(0..9)) as ArrayRef,
+ Arc::new(Float32Array::from(vec![
+ Some(float_negative_nan),
+ Some(f32::from_bits(0xffff_1234)),
+ Some(f32::NAN),
+ Some(f32::from_bits(0x7fc0_0010)),
+ Some(-1.0),
+ Some(-0.0),
+ Some(0.0),
+ Some(1.0),
+ None,
+ ])) as ArrayRef,
+ Arc::new(Float64Array::from(vec![
+ Some(double_negative_nan),
+ Some(f64::from_bits(0xffff_1234_5678_9abc)),
+ Some(f64::NAN),
+ Some(f64::from_bits(0x7ff8_0000_0000_0010)),
+ Some(-1.0),
+ Some(-0.0),
+ Some(0.0),
+ Some(1.0),
+ None,
+ ])) as ArrayRef,
+ ],
+ )
+ .unwrap();
+ let mut table_write = TableWrite::new(&table,
"test-user".to_string()).unwrap();
+ table_write.write_arrow_batch(&batch).await.unwrap();
+ let messages = table_write.prepare_commit().await.unwrap();
+ TableCommit::new(table.clone(), "test-user".to_string())
+ .commit(messages)
+ .await
+ .unwrap();
+
+ for column in ["f", "d"] {
+ let shard_count = table
+ .new_btree_global_index_build_builder()
+ .with_index_column(column)
+ .with_index_type(BITMAP_GLOBAL_INDEX_TYPE)
+ .execute()
+ .await
+ .unwrap();
+ assert_eq!(shard_count, 1);
+ }
+
+ let mut disabled_options = table.schema().options().clone();
+ disabled_options.insert("global-index.enabled".to_string(),
"false".to_string());
+ let table_without_index = Table::new(
+ table.file_io().clone(),
+ table.identifier().clone(),
+ table.location().to_string(),
+ table.schema().copy_with_replaced_options(disabled_options),
+ None,
+ );
+
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ let cases = [
+ (
+ "Float < 0",
+ predicates.less_than("f", Datum::Float(0.0)).unwrap(),
+ vec![0, 1, 4, 5],
+ ),
+ (
+ "Double < 0",
+ predicates.less_than("d", Datum::Double(0.0)).unwrap(),
+ vec![0, 1, 4, 5],
+ ),
+ (
+ "Float = canonical NaN",
+ predicates.equal("f", Datum::Float(f32::NAN)).unwrap(),
+ vec![2],
+ ),
+ (
+ "Double = canonical NaN",
+ predicates.equal("d", Datum::Double(f64::NAN)).unwrap(),
+ vec![2],
+ ),
+ (
+ "Float = negative NaN",
+ predicates
+ .equal("f", Datum::Float(float_negative_nan))
+ .unwrap(),
+ vec![0],
+ ),
+ (
+ "Double = negative NaN",
+ predicates
+ .equal("d", Datum::Double(double_negative_nan))
+ .unwrap(),
+ vec![0],
+ ),
+ (
+ "Float IN NaNs",
+ predicates
+ .is_in(
+ "f",
+ vec![Datum::Float(float_negative_nan),
Datum::Float(f32::NAN)],
+ )
+ .unwrap(),
+ vec![0, 2],
+ ),
+ (
+ "Double IN NaNs",
+ predicates
+ .is_in(
+ "d",
+ vec![Datum::Double(double_negative_nan),
Datum::Double(f64::NAN)],
+ )
+ .unwrap(),
+ vec![0, 2],
+ ),
+ (
+ "Float != canonical NaN",
+ predicates.not_equal("f", Datum::Float(f32::NAN)).unwrap(),
+ vec![0, 1, 3, 4, 5, 6, 7],
+ ),
+ (
+ "Double != canonical NaN",
+ predicates.not_equal("d", Datum::Double(f64::NAN)).unwrap(),
+ vec![0, 1, 3, 4, 5, 6, 7],
+ ),
+ (
+ "Float NOT IN",
+ predicates
+ .is_not_in("f", vec![Datum::Float(f32::NAN),
Datum::Float(0.0)])
+ .unwrap(),
+ vec![0, 1, 3, 4, 5, 7],
+ ),
+ (
+ "Double NOT IN",
+ predicates
+ .is_not_in("d", vec![Datum::Double(f64::NAN),
Datum::Double(0.0)])
+ .unwrap(),
+ vec![0, 1, 3, 4, 5, 7],
+ ),
+ (
+ "Float combined range",
+ Predicate::and(vec![
+ predicates
+ .greater_or_equal("f",
Datum::Float(float_negative_nan))
+ .unwrap(),
+ predicates.less_or_equal("f", Datum::Float(0.0)).unwrap(),
+ ]),
+ vec![0, 4, 5, 6],
+ ),
+ (
+ "Double combined range",
+ Predicate::and(vec![
+ predicates
+ .greater_or_equal("d",
Datum::Double(double_negative_nan))
+ .unwrap(),
+ predicates.less_or_equal("d", Datum::Double(0.0)).unwrap(),
+ ]),
+ vec![0, 4, 5, 6],
+ ),
+ ];
+
+ for (name, predicate, expected) in cases {
+ let without_index = scan_ids(&table_without_index,
predicate.clone()).await;
+ assert_eq!(without_index, expected, "{name}: residual baseline");
+ let with_index = scan_ids(&table, predicate).await;
+ assert_eq!(
+ with_index, without_index,
+ "{name}: global index changed rows"
+ );
+ }
+ }
+
/// Bitmap is built through the same sorted builder; a second build with no
/// new data must be a no-op keyed on the bitmap coverage — not error, and
/// not be confused by any btree coverage of the same field.
diff --git a/crates/paimon/src/table/global_index_scanner.rs
b/crates/paimon/src/table/global_index_scanner.rs
index c0122d44..ad0fa022 100644
--- a/crates/paimon/src/table/global_index_scanner.rs
+++ b/crates/paimon/src/table/global_index_scanner.rs
@@ -20,7 +20,10 @@
//!
//! Reference:
[org.apache.paimon.index.GlobalIndexScanner](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/index/GlobalIndexScanner.java)
-use super::bitmap_global_index_reader::BitmapGlobalIndexReader;
+use super::bitmap_global_index_reader::{
+ is_bitmap_floating_residual_sensitive_op, make_bitmap_key_comparator,
serialize_bitmap_datum,
+ BitmapGlobalIndexReader,
+};
use super::global_index_types::{
normalize_sorted_global_index_type, BITMAP_GLOBAL_INDEX_TYPE,
BTREE_GLOBAL_INDEX_TYPE,
};
@@ -113,6 +116,40 @@ enum GlobalIndexFileKind {
Bitmap,
}
+fn is_floating_point(data_type: &DataType) -> bool {
+ matches!(data_type, DataType::Float(_) | DataType::Double(_))
+}
+
+fn bitmap_meta_may_match(
+ meta: &BTreeIndexMeta,
+ op: PredicateOperator,
+ data_type: &DataType,
+ serialized_literals: &[Vec<u8>],
+ cmp: &dyn Fn(&[u8], &[u8]) -> Ordering,
+) -> bool {
+ if is_floating_point(data_type) &&
is_bitmap_floating_residual_sensitive_op(op) {
+ !meta.only_nulls()
+ } else {
+ meta.may_match(op, serialized_literals, cmp)
+ }
+}
+
+fn bitmap_meta_may_match_between(
+ meta: &BTreeIndexMeta,
+ data_type: &DataType,
+ from_key: &[u8],
+ to_key: &[u8],
+ cmp: &dyn Fn(&[u8], &[u8]) -> Ordering,
+) -> bool {
+ if is_floating_point(data_type)
+ && is_bitmap_floating_residual_sensitive_op(PredicateOperator::Between)
+ {
+ !meta.only_nulls()
+ } else {
+ meta.may_match_between(from_key, to_key, cmp)
+ }
+}
+
impl GlobalIndexFileKind {
fn name(self) -> &'static str {
match self {
@@ -444,23 +481,48 @@ impl GlobalIndexScanner {
let pruning_info: Vec<_> = effective_predicates
.iter()
.map(|(op, literals, data_type)| {
- let cmp = make_key_comparator(data_type);
- let serialized: Vec<Vec<u8>> = literals
+ let btree_cmp = make_key_comparator(data_type);
+ let btree_serialized = literals
.iter()
.map(|l| serialize_datum(l, data_type))
- .collect();
- (*op, cmp, serialized)
+ .collect::<Vec<_>>();
+ let bitmap_cmp = make_bitmap_key_comparator(data_type);
+ let bitmap_serialized = literals
+ .iter()
+ .map(|l| serialize_bitmap_datum(l, data_type))
+ .collect::<Vec<_>>();
+ (
+ *op,
+ *data_type,
+ btree_cmp,
+ btree_serialized,
+ bitmap_cmp,
+ bitmap_serialized,
+ )
})
.collect();
let predicate_matches: Vec<Vec<bool>> = pruning_info
.iter()
- .map(|(op, cmp, serialized)| {
- entries
- .iter()
- .map(|entry| entry.meta.may_match(*op, serialized, cmp))
- .collect()
- })
+ .map(
+ |(op, data_type, btree_cmp, btree_serialized, bitmap_cmp,
bitmap_serialized)| {
+ entries
+ .iter()
+ .map(|entry| match entry.index_type {
+ GlobalIndexFileKind::BTree => {
+ entry.meta.may_match(*op, btree_serialized,
btree_cmp)
+ }
+ GlobalIndexFileKind::Bitmap =>
bitmap_meta_may_match(
+ &entry.meta,
+ *op,
+ data_type,
+ bitmap_serialized,
+ bitmap_cmp.as_ref(),
+ ),
+ })
+ .collect()
+ },
+ )
.collect();
let predicate_fallback_plans: Vec<Option<FallbackScanPlan>> =
effective_predicates
.iter()
@@ -473,12 +535,28 @@ impl GlobalIndexScanner {
let between_matches_by_entry: Vec<bool> = match between.as_ref() {
Some(b) => {
- let cmp = make_key_comparator(b.data_type);
- let from_key = serialize_datum(b.from, b.data_type);
- let to_key = serialize_datum(b.to, b.data_type);
+ let btree_cmp = make_key_comparator(b.data_type);
+ let btree_from = serialize_datum(b.from, b.data_type);
+ let btree_to = serialize_datum(b.to, b.data_type);
+ let bitmap_cmp = make_bitmap_key_comparator(b.data_type);
+ let bitmap_from = serialize_bitmap_datum(b.from, b.data_type);
+ let bitmap_to = serialize_bitmap_datum(b.to, b.data_type);
entries
.iter()
- .map(|entry| entry.meta.may_match_between(&from_key,
&to_key, &cmp))
+ .map(|entry| match entry.index_type {
+ GlobalIndexFileKind::BTree => {
+ entry
+ .meta
+ .may_match_between(&btree_from, &btree_to,
&btree_cmp)
+ }
+ GlobalIndexFileKind::Bitmap =>
bitmap_meta_may_match_between(
+ &entry.meta,
+ b.data_type,
+ &bitmap_from,
+ &bitmap_to,
+ bitmap_cmp.as_ref(),
+ ),
+ })
.collect()
}
None => Vec::new(),
@@ -608,8 +686,12 @@ impl GlobalIndexScanner {
if plan.between_matches && plan.between_evaluated {
let between = between.expect("evaluated between query is present");
- let from_key = serialize_datum(between.from, between.data_type);
- let to_key = serialize_datum(between.to, between.data_type);
+ let serialize_key = match entry.index_type {
+ GlobalIndexFileKind::BTree => serialize_datum,
+ GlobalIndexFileKind::Bitmap => serialize_bitmap_datum,
+ };
+ let from_key = serialize_key(between.from, between.data_type);
+ let to_key = serialize_key(between.to, between.data_type);
let bitmap = reader
.as_ref()
.expect("reader is opened when between matches")
@@ -1332,6 +1414,9 @@ pub(crate) async fn evaluate_global_index(
#[cfg(test)]
mod tests {
use super::*;
+ use crate::btree::test_util::VecFileWrite;
+ use crate::btree::{BTreeIndexWriter, BlockCompressionType};
+ use crate::table::bitmap_global_index_reader::BitmapGlobalIndexWriter;
use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
use std::sync::Arc;
@@ -1419,6 +1504,129 @@ mod tests {
assert_eq!(key, b"hello".to_vec());
}
+ fn assert_bitmap_floating_meta_policy(
+ data_type: DataType,
+ min: Datum,
+ max: Datum,
+ outside: Datum,
+ nan: Datum,
+ ) {
+ let cmp = make_bitmap_key_comparator(&data_type);
+ let min_key = serialize_bitmap_datum(&min, &data_type);
+ let max_key = serialize_bitmap_datum(&max, &data_type);
+ let outside_key = serialize_bitmap_datum(&outside, &data_type);
+ let nan_key = serialize_bitmap_datum(&nan, &data_type);
+ let meta = BTreeIndexMeta::new(Some(min_key.clone()), Some(max_key),
false);
+
+ assert!(!bitmap_meta_may_match(
+ &meta,
+ PredicateOperator::Eq,
+ &data_type,
+ std::slice::from_ref(&outside_key),
+ cmp.as_ref(),
+ ));
+ assert!(!bitmap_meta_may_match(
+ &meta,
+ PredicateOperator::In,
+ &data_type,
+ std::slice::from_ref(&outside_key),
+ cmp.as_ref(),
+ ));
+ assert!(!bitmap_meta_may_match(
+ &meta,
+ PredicateOperator::IsNull,
+ &data_type,
+ &[],
+ cmp.as_ref(),
+ ));
+ assert!(bitmap_meta_may_match(
+ &meta,
+ PredicateOperator::IsNotNull,
+ &data_type,
+ &[],
+ cmp.as_ref(),
+ ));
+
+ let nan_meta = BTreeIndexMeta::new(Some(min_key),
Some(nan_key.clone()), false);
+ assert!(bitmap_meta_may_match(
+ &nan_meta,
+ PredicateOperator::Eq,
+ &data_type,
+ std::slice::from_ref(&nan_key),
+ cmp.as_ref(),
+ ));
+ assert!(bitmap_meta_may_match(
+ &nan_meta,
+ PredicateOperator::In,
+ &data_type,
+ std::slice::from_ref(&nan_key),
+ cmp.as_ref(),
+ ));
+
+ assert!(bitmap_meta_may_match(
+ &meta,
+ PredicateOperator::Gt,
+ &data_type,
+ std::slice::from_ref(&outside_key),
+ cmp.as_ref(),
+ ));
+ assert!(bitmap_meta_may_match_between(
+ &meta,
+ &data_type,
+ &outside_key,
+ &outside_key,
+ cmp.as_ref(),
+ ));
+
+ let only_nulls = BTreeIndexMeta::new(None, None, true);
+ assert!(bitmap_meta_may_match(
+ &only_nulls,
+ PredicateOperator::IsNull,
+ &data_type,
+ &[],
+ cmp.as_ref(),
+ ));
+ assert!(!bitmap_meta_may_match(
+ &only_nulls,
+ PredicateOperator::IsNotNull,
+ &data_type,
+ &[],
+ cmp.as_ref(),
+ ));
+ assert!(!bitmap_meta_may_match(
+ &only_nulls,
+ PredicateOperator::NotEq,
+ &data_type,
+ std::slice::from_ref(&outside_key),
+ cmp.as_ref(),
+ ));
+ assert!(!bitmap_meta_may_match_between(
+ &only_nulls,
+ &data_type,
+ &outside_key,
+ &outside_key,
+ cmp.as_ref(),
+ ));
+ }
+
+ #[test]
+ fn test_bitmap_floating_meta_prunes_equality_and_fails_open_for_ranges() {
+ assert_bitmap_floating_meta_policy(
+ DataType::Float(crate::spec::FloatType::new()),
+ Datum::Float(-1.0),
+ Datum::Float(1.0),
+ Datum::Float(2.0),
+ Datum::Float(f32::NAN),
+ );
+ assert_bitmap_floating_meta_policy(
+ DataType::Double(crate::spec::DoubleType::new()),
+ Datum::Double(-1.0),
+ Datum::Double(1.0),
+ Datum::Double(2.0),
+ Datum::Double(f64::NAN),
+ );
+ }
+
#[test]
fn test_row_range_index_merges_overlapping() {
let idx = RowRangeIndex::create(vec![
@@ -1517,24 +1725,27 @@ mod tests {
(file_io, table_path, testdata_name.to_string(), tmp)
}
- type JavaBitmapTestdataTable = (FileIO, String, String, BTreeIndexMeta,
tempfile::TempDir);
+ type BitmapTestdataTable = (FileIO, String, String, BTreeIndexMeta,
tempfile::TempDir);
- fn setup_java_bitmap_testdata_table() -> JavaBitmapTestdataTable {
- const FILE_NAME: &str = "bitmap_varchar_java.index";
- let src = format!("{}/testdata/bitmap/{FILE_NAME}",
env!("CARGO_MANIFEST_DIR"));
+ fn setup_bitmap_testdata_table(file_name: &str) -> BitmapTestdataTable {
+ let src = format!("{}/testdata/bitmap/{file_name}",
env!("CARGO_MANIFEST_DIR"));
let meta_src = format!(
- "{}/testdata/bitmap/{FILE_NAME}.meta",
+ "{}/testdata/bitmap/{file_name}.meta",
env!("CARGO_MANIFEST_DIR")
);
let tmp = tempfile::tempdir().unwrap();
let index_dir = tmp.path().join("index");
std::fs::create_dir_all(&index_dir).unwrap();
- std::fs::copy(&src, index_dir.join(FILE_NAME)).unwrap();
+ std::fs::copy(&src, index_dir.join(file_name)).unwrap();
let meta =
BTreeIndexMeta::deserialize(&std::fs::read(meta_src).unwrap()).unwrap();
let table_path = format!("file://{}", tmp.path().display());
let file_io = crate::io::FileIOBuilder::new("file").build().unwrap();
- (file_io, table_path, FILE_NAME.to_string(), meta, tmp)
+ (file_io, table_path, file_name.to_string(), meta, tmp)
+ }
+
+ fn setup_java_bitmap_testdata_table() -> BitmapTestdataTable {
+ setup_bitmap_testdata_table("bitmap_varchar_java.index")
}
fn make_global_index_entry(
@@ -2054,6 +2265,317 @@ mod tests {
assert_eq!(null_result.unwrap(), vec![RowRange::new(104, 104)]);
}
+ async fn assert_bitmap_int_fixture(file_name: &str) {
+ let data_type = DataType::Int(crate::spec::IntType::new());
+ let (file_io, table_path, file_name, meta, _tmp) =
setup_bitmap_testdata_table(file_name);
+ let entries = vec![make_global_index_entry_with_type(
+ BITMAP_GLOBAL_INDEX_TYPE,
+ &file_name,
+ 1,
+ 100,
+ 105,
+ &meta,
+ )];
+ let fields = int_schema_fields();
+ assert_eq!(meta.first_key, Some(le_int_key(-1)));
+ assert_eq!(meta.last_key, Some(le_int_key(256)));
+ assert!(meta.has_nulls);
+
+ let cases = [
+ (
+ PredicateOperator::Eq,
+ vec![Datum::Int(0)],
+ vec![RowRange::new(101, 102)],
+ ),
+ (
+ PredicateOperator::Eq,
+ vec![Datum::Int(256)],
+ vec![RowRange::new(104, 104)],
+ ),
+ (
+ PredicateOperator::In,
+ vec![Datum::Int(-1), Datum::Int(1), Datum::Int(256)],
+ vec![RowRange::new(100, 100), RowRange::new(103, 104)],
+ ),
+ (
+ PredicateOperator::NotEq,
+ vec![Datum::Int(0)],
+ vec![RowRange::new(100, 100), RowRange::new(103, 104)],
+ ),
+ (
+ PredicateOperator::NotIn,
+ vec![Datum::Int(-1), Datum::Int(1), Datum::Int(256)],
+ vec![RowRange::new(101, 102)],
+ ),
+ (
+ PredicateOperator::IsNull,
+ vec![],
+ vec![RowRange::new(105, 105)],
+ ),
+ ];
+
+ for (op, literals, expected) in cases {
+ let predicates = vec![Predicate::Leaf {
+ column: "id".to_string(),
+ index: 0,
+ data_type: data_type.clone(),
+ op,
+ literals,
+ }];
+ let result =
+ evaluate_global_index_fast(&file_io, &table_path, &entries,
&predicates, &fields)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(result, expected, "{file_name}: {op}");
+ }
+ }
+
+ #[tokio::test]
+ async fn test_evaluate_java_logical_order_bitmap_int_fixture() {
+ assert_bitmap_int_fixture("bitmap_int_logical_java.index").await;
+ }
+
+ async fn assert_bitmap_nan_equality_uses_direct_lookup(
+ data_type: DataType,
+ nan_literals: [Datum; 3],
+ zero: Datum,
+ ) {
+ let output = VecFileWrite::new();
+ let captured = output.clone();
+ let mut writer = BitmapGlobalIndexWriter::new(
+ Box::new(output),
+ 1,
+ BlockCompressionType::None,
+ make_bitmap_key_comparator(&data_type),
+ );
+ for (row_id, literal) in nan_literals.iter().enumerate() {
+ let key = serialize_bitmap_datum(literal, &data_type);
+ writer.write(Some(&key), row_id as i64).unwrap();
+ }
+ let zero_key = serialize_bitmap_datum(&zero, &data_type);
+ writer.write(Some(&zero_key), 3).unwrap();
+ let write_result = writer.finish().await.unwrap();
+ let bytes = captured.to_vec();
+
+ let tmp = tempfile::tempdir().unwrap();
+ let index_dir = tmp.path().join("index");
+ std::fs::create_dir_all(&index_dir).unwrap();
+ let file_name = "bitmap-current.index";
+ std::fs::write(index_dir.join(file_name), &bytes).unwrap();
+ let table_path = format!("file://{}", tmp.path().display());
+ let file_io = crate::io::FileIOBuilder::new("file").build().unwrap();
+
+ let mut entry = make_global_index_entry_with_type(
+ BITMAP_GLOBAL_INDEX_TYPE,
+ file_name,
+ 1,
+ 100,
+ 103,
+ &write_result.meta,
+ );
+ entry.index_file.file_size = bytes.len() as i64;
+ let entries = vec![entry];
+ let fields = vec![DataField::new(1, "id".to_string(),
data_type.clone())];
+ let cases = [
+ (PredicateOperator::Eq, vec![nan_literals[0].clone()]),
+ (
+ PredicateOperator::In,
+ vec![nan_literals[1].clone(), nan_literals[2].clone()],
+ ),
+ ];
+
+ for (op, literals) in cases {
+ let predicates = vec![Predicate::Leaf {
+ column: "id".to_string(),
+ index: 0,
+ data_type: data_type.clone(),
+ op,
+ literals,
+ }];
+ let result = evaluate_global_index_fast_with_fallback_size(
+ &file_io,
+ &table_path,
+ &entries,
+ &predicates,
+ &fields,
+ i64::MAX,
+ 0,
+ )
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(result, vec![RowRange::new(100, 102)], "{data_type:?}:
{op}");
+ }
+ }
+
+ #[tokio::test]
+ async fn
test_bitmap_nan_equality_uses_direct_lookup_with_fallback_scan_disabled() {
+ assert_bitmap_nan_equality_uses_direct_lookup(
+ DataType::Float(crate::spec::FloatType::new()),
+ [
+ Datum::Float(f32::from_bits(0xffc0_0001)),
+ Datum::Float(f32::from_bits(0x7fc0_0010)),
+ Datum::Float(f32::NAN),
+ ],
+ Datum::Float(0.0),
+ )
+ .await;
+ assert_bitmap_nan_equality_uses_direct_lookup(
+ DataType::Double(crate::spec::DoubleType::new()),
+ [
+ Datum::Double(f64::from_bits(0xfff8_0000_0000_0001)),
+ Datum::Double(f64::from_bits(0x7ff8_0000_0000_0010)),
+ Datum::Double(f64::NAN),
+ ],
+ Datum::Double(0.0),
+ )
+ .await;
+ }
+
+ fn legacy_floating_comparator(data_type: &DataType) -> BoxedCmp {
+ match data_type {
+ DataType::Float(_) => Box::new(|left, right| {
+ let left = f32::from_le_bytes(left.try_into().unwrap());
+ let right = f32::from_le_bytes(right.try_into().unwrap());
+ left.total_cmp(&right)
+ }),
+ DataType::Double(_) => Box::new(|left, right| {
+ let left = f64::from_le_bytes(left.try_into().unwrap());
+ let right = f64::from_le_bytes(right.try_into().unwrap());
+ left.total_cmp(&right)
+ }),
+ _ => unreachable!("legacy floating comparator requires Float or
Double"),
+ }
+ }
+
+ async fn assert_legacy_floating_btree(
+ file_name: &str,
+ data_type: DataType,
+ nan_keys: Vec<Vec<u8>>,
+ nan_literals: Vec<Datum>,
+ zero_key: Vec<u8>,
+ zero_literal: Datum,
+ ) {
+ let mut rows = nan_keys
+ .into_iter()
+ .enumerate()
+ .map(|(row_id, key)| (key, row_id as i64))
+ .collect::<Vec<_>>();
+ rows.push((zero_key, 3));
+ let cmp = legacy_floating_comparator(&data_type);
+ rows.sort_by(|left, right| cmp(&left.0, &right.0));
+ let expected_first_key = rows.first().unwrap().0.clone();
+ let expected_last_key = rows.last().unwrap().0.clone();
+
+ let output = VecFileWrite::new();
+ let captured = output.clone();
+ let mut writer =
+ BTreeIndexWriter::with_comparator(Box::new(output), 1,
BlockCompressionType::None, cmp);
+ for (key, row_id) in rows {
+ writer.write(Some(&key), row_id).await.unwrap();
+ }
+ let write_result = writer.finish().await.unwrap();
+ assert_eq!(write_result.meta.first_key, Some(expected_first_key));
+ assert_eq!(write_result.meta.last_key, Some(expected_last_key));
+
+ let tmp = tempfile::tempdir().unwrap();
+ let index_dir = tmp.path().join("index");
+ std::fs::create_dir_all(&index_dir).unwrap();
+ std::fs::write(index_dir.join(file_name), captured.to_vec()).unwrap();
+ let table_path = format!("file://{}", tmp.path().display());
+ let file_io = crate::io::FileIOBuilder::new("file").build().unwrap();
+ let entries = vec![make_global_index_entry(
+ file_name,
+ 1,
+ 100,
+ 103,
+ &write_result.meta,
+ )];
+ let fields = vec![DataField::new(1, "id".to_string(),
data_type.clone())];
+ let cases = [
+ (
+ PredicateOperator::Eq,
+ vec![zero_literal.clone()],
+ vec![RowRange::new(103, 103)],
+ ),
+ (
+ PredicateOperator::Eq,
+ vec![nan_literals[0].clone()],
+ vec![RowRange::new(100, 100)],
+ ),
+ (
+ PredicateOperator::In,
+ vec![
+ nan_literals[0].clone(),
+ nan_literals[1].clone(),
+ zero_literal,
+ ],
+ vec![RowRange::new(100, 101), RowRange::new(103, 103)],
+ ),
+ ];
+
+ for (op, literals, expected) in cases {
+ let predicates = vec![Predicate::Leaf {
+ column: "id".to_string(),
+ index: 0,
+ data_type: data_type.clone(),
+ op,
+ literals,
+ }];
+ let result =
+ evaluate_global_index_fast(&file_io, &table_path, &entries,
&predicates, &fields)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(result, expected, "{file_name}: {op}");
+ }
+ }
+
+ #[tokio::test]
+ async fn test_evaluate_legacy_float_btree() {
+ let nan_bits = [0xffc0_0001u32, 0xffc0_0010, 0xffff_1234];
+ assert_legacy_floating_btree(
+ "btree_float_legacy_rust.index",
+ DataType::Float(crate::spec::FloatType::new()),
+ nan_bits
+ .iter()
+ .map(|bits| bits.to_le_bytes().to_vec())
+ .collect(),
+ nan_bits
+ .iter()
+ .map(|bits| Datum::Float(f32::from_bits(*bits)))
+ .collect(),
+ 0.0f32.to_le_bytes().to_vec(),
+ Datum::Float(0.0),
+ )
+ .await;
+ }
+
+ #[tokio::test]
+ async fn test_evaluate_legacy_double_btree() {
+ let nan_bits = [
+ 0xfff8_0000_0000_0001u64,
+ 0xfff8_0000_0000_0010,
+ 0xffff_1234_5678_9abc,
+ ];
+ assert_legacy_floating_btree(
+ "btree_double_legacy_rust.index",
+ DataType::Double(crate::spec::DoubleType::new()),
+ nan_bits
+ .iter()
+ .map(|bits| bits.to_le_bytes().to_vec())
+ .collect(),
+ nan_bits
+ .iter()
+ .map(|bits| Datum::Double(f64::from_bits(*bits)))
+ .collect(),
+ 0.0f64.to_le_bytes().to_vec(),
+ Datum::Double(0.0),
+ )
+ .await;
+ }
+
#[tokio::test]
async fn test_evaluate_java_bitmap_golden_index_string_fallback_scan() {
let data_type =
DataType::VarChar(crate::spec::VarCharType::string_type());
diff --git a/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index
b/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index
new file mode 100644
index 00000000..91abd954
Binary files /dev/null and
b/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index differ
diff --git a/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index.meta
b/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index.meta
new file mode 100644
index 00000000..306c6cb7
Binary files /dev/null and
b/crates/paimon/testdata/bitmap/bitmap_int_logical_java.index.meta differ