mirror of
https://github.com/quickwit-oss/tantivy.git
synced 2026-08-18 03:58:21 +00:00
issue/55 working.
This commit is contained in:
+30
-27
@@ -8,19 +8,17 @@ pub fn compute_num_bits(amplitude: u32) -> u8 {
|
||||
(32u32 - amplitude.leading_zeros()) as u8
|
||||
}
|
||||
|
||||
pub struct BitPacker<TWrite: Write> {
|
||||
output: TWrite,
|
||||
pub struct BitPacker {
|
||||
mini_buffer: u64,
|
||||
mini_buffer_written: usize,
|
||||
num_bits: usize,
|
||||
written_size: usize,
|
||||
}
|
||||
|
||||
impl<TWrite: Write> BitPacker<TWrite> {
|
||||
impl BitPacker {
|
||||
|
||||
pub fn new(output: TWrite, num_bits: usize) -> BitPacker<TWrite> {
|
||||
pub fn new(num_bits: usize) -> BitPacker {
|
||||
BitPacker {
|
||||
output: output,
|
||||
mini_buffer: 0u64,
|
||||
mini_buffer_written: 0,
|
||||
num_bits: num_bits,
|
||||
@@ -28,11 +26,11 @@ impl<TWrite: Write> BitPacker<TWrite> {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn write(&mut self, val: u32) -> io::Result<()> {
|
||||
pub fn write<TWrite: Write>(&mut self, val: u32, output: &mut TWrite) -> io::Result<()> {
|
||||
let val_u64 = val as u64;
|
||||
if self.mini_buffer_written + self.num_bits > 64 {
|
||||
self.mini_buffer |= val_u64.wrapping_shl(self.mini_buffer_written as u32);
|
||||
self.written_size += self.mini_buffer.serialize(&mut self.output)?;
|
||||
self.written_size += self.mini_buffer.serialize(output)?;
|
||||
self.mini_buffer = val_u64.wrapping_shr((64 - self.mini_buffer_written) as u32);
|
||||
self.mini_buffer_written = self.mini_buffer_written + (self.num_bits as usize) - 64;
|
||||
}
|
||||
@@ -40,7 +38,7 @@ impl<TWrite: Write> BitPacker<TWrite> {
|
||||
self.mini_buffer |= val_u64 << self.mini_buffer_written;
|
||||
self.mini_buffer_written += self.num_bits;
|
||||
if self.mini_buffer_written == 64 {
|
||||
self.written_size += self.mini_buffer.serialize(&mut self.output)?;
|
||||
self.written_size += self.mini_buffer.serialize(output)?;
|
||||
self.mini_buffer_written = 0;
|
||||
self.mini_buffer = 0u64;
|
||||
}
|
||||
@@ -48,37 +46,39 @@ impl<TWrite: Write> BitPacker<TWrite> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn flush(&mut self) -> io::Result<()>{
|
||||
fn flush<TWrite: Write>(&mut self, output: &mut TWrite) -> io::Result<()>{
|
||||
if self.mini_buffer_written > 0 {
|
||||
let num_bytes = (self.mini_buffer_written + 7) / 8;
|
||||
let arr: [u8; 8] = unsafe { mem::transmute::<u64, [u8; 8]>(self.mini_buffer) };
|
||||
self.output.write_all(&arr[..num_bytes])?;
|
||||
output.write_all(&arr[..num_bytes])?;
|
||||
self.written_size += num_bytes;
|
||||
self.mini_buffer_written = 0;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn close(mut self) -> io::Result<(TWrite, usize)> {
|
||||
self.flush()?;
|
||||
Ok((self.output, self.written_size))
|
||||
pub fn close<TWrite: Write>(&mut self, output: &mut TWrite) -> io::Result<usize> {
|
||||
self.flush(output)?;
|
||||
Ok(self.written_size)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
pub struct BitUnpacker<'a> {
|
||||
data: &'a [u8],
|
||||
pub struct BitUnpacker {
|
||||
num_bits: usize,
|
||||
mask: u32,
|
||||
data_ptr: *const u8,
|
||||
data_len: usize,
|
||||
}
|
||||
|
||||
impl<'a> BitUnpacker<'a> {
|
||||
pub fn new(data: &'a [u8], num_bits: usize) -> BitUnpacker<'a> {
|
||||
impl BitUnpacker {
|
||||
pub fn new(data: &[u8], num_bits: usize) -> BitUnpacker {
|
||||
BitUnpacker {
|
||||
data: data,
|
||||
num_bits: num_bits,
|
||||
mask: (1u32 << num_bits) - 1u32,
|
||||
data_ptr: data.as_ptr(),
|
||||
data_len: data.len()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -87,15 +87,17 @@ impl<'a> BitUnpacker<'a> {
|
||||
return 0;
|
||||
}
|
||||
let addr = (idx * self.num_bits) / 8;
|
||||
let bit_shift = (idx * self.num_bits) - addr * 8;
|
||||
let bit_shift = idx * self.num_bits - addr * 8;
|
||||
let val_unshifted_unmasked: u64;
|
||||
if addr + 8 <= self.data.len() {
|
||||
val_unshifted_unmasked = unsafe { * (self.data.as_ptr().offset(addr as isize) as *const u64) };
|
||||
if addr + 8 <= self.data_len {
|
||||
val_unshifted_unmasked = unsafe { * (self.data_ptr.offset(addr as isize) as *const u64) };
|
||||
}
|
||||
else {
|
||||
else {
|
||||
let mut arr = [0u8; 8];
|
||||
for i in 0..self.data.len() - addr {
|
||||
arr[i] = self.data[addr + i];
|
||||
if addr < self.data_len {
|
||||
for i in 0..self.data_len - addr {
|
||||
arr[i] = unsafe { *self.data_ptr.offset( (addr + i) as isize) };
|
||||
}
|
||||
}
|
||||
val_unshifted_unmasked = unsafe { mem::transmute::<[u8; 8], u64>(arr) };
|
||||
}
|
||||
@@ -124,7 +126,8 @@ mod test {
|
||||
}
|
||||
|
||||
fn test_bitpacker_util(len: usize, num_bits: usize) {
|
||||
let mut bitpacker = BitPacker::new(Vec::new(), num_bits);
|
||||
let mut data = Vec::new();
|
||||
let mut bitpacker = BitPacker::new(num_bits);
|
||||
let max_val: u32 = (1 << num_bits) - 1;
|
||||
let vals: Vec<u32> = (0u32..len as u32).map(|i| {
|
||||
if max_val == 0 {
|
||||
@@ -135,9 +138,9 @@ mod test {
|
||||
}
|
||||
}).collect();
|
||||
for &val in &vals {
|
||||
bitpacker.write(val).unwrap();
|
||||
bitpacker.write(val, &mut data).unwrap();
|
||||
}
|
||||
let (data, num_bytes) = bitpacker.close().unwrap();
|
||||
let num_bytes = bitpacker.close(&mut data).unwrap();
|
||||
assert_eq!(num_bytes, (num_bits * len + 7) / 8);
|
||||
assert_eq!(data.len(), num_bytes);
|
||||
let bitunpacker = BitUnpacker::new(&data, num_bits);
|
||||
|
||||
@@ -71,12 +71,11 @@ mod compression {
|
||||
}
|
||||
let num_bits = compute_num_bits(max_delta);
|
||||
output.write_all(&[num_bits]).unwrap();
|
||||
let mut bit_packer = BitPacker::new(output, num_bits as usize);
|
||||
let mut bit_packer = BitPacker::new(num_bits as usize);
|
||||
for val in &deltas {
|
||||
bit_packer.write(*val).unwrap();
|
||||
bit_packer.write(*val, &mut output).unwrap();
|
||||
}
|
||||
let (_, written_size) = bit_packer.close().expect("packing in memory should never fail");
|
||||
written_size + 1
|
||||
1 + bit_packer.close(&mut output).expect("packing in memory should never fail")
|
||||
}
|
||||
|
||||
pub fn uncompress_sorted(compressed_data: &[u8], output: &mut [u32], mut offset: u32) -> usize {
|
||||
@@ -95,12 +94,11 @@ mod compression {
|
||||
let max = vals.iter().cloned().max().expect("compress unsorted called with an empty array");
|
||||
let num_bits = compute_num_bits(max);
|
||||
output.write_all(&[num_bits]).unwrap();
|
||||
let mut bit_packer = BitPacker::new(output, num_bits as usize);
|
||||
let mut bit_packer = BitPacker::new(num_bits as usize);
|
||||
for val in vals {
|
||||
bit_packer.write(*val).unwrap();
|
||||
bit_packer.write(*val, &mut output).unwrap();
|
||||
}
|
||||
let (_, written_size) = bit_packer.close().expect("packing in memory should never fail");
|
||||
1 + written_size
|
||||
1 + bit_packer.close(&mut output).expect("packing in memory should never fail")
|
||||
}
|
||||
|
||||
pub fn uncompress_unsorted(compressed_data: &[u8], output: &mut [u32]) -> usize {
|
||||
|
||||
@@ -18,8 +18,6 @@ pub use self::writer::{U32FastFieldsWriter, U32FastFieldWriter};
|
||||
pub use self::reader::{U32FastFieldsReader, U32FastFieldReader};
|
||||
pub use self::serializer::FastFieldSerializer;
|
||||
|
||||
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -79,7 +77,7 @@ mod tests {
|
||||
}
|
||||
let source = directory.open_read(&path).unwrap();
|
||||
{
|
||||
assert_eq!(source.len(), 26 as usize);
|
||||
assert_eq!(source.len(), 20 as usize);
|
||||
}
|
||||
{
|
||||
let fast_field_readers = U32FastFieldsReader::open(source).unwrap();
|
||||
@@ -112,7 +110,7 @@ mod tests {
|
||||
}
|
||||
let source = directory.open_read(&path).unwrap();
|
||||
{
|
||||
assert_eq!(source.len(), 50 as usize);
|
||||
assert_eq!(source.len(), 45 as usize);
|
||||
}
|
||||
{
|
||||
let fast_field_readers = U32FastFieldsReader::open(source).unwrap();
|
||||
@@ -147,7 +145,7 @@ mod tests {
|
||||
}
|
||||
let source = directory.open_read(&path).unwrap();
|
||||
{
|
||||
assert_eq!(source.len(), 26 as usize);
|
||||
assert_eq!(source.len(), 18 as usize);
|
||||
}
|
||||
{
|
||||
let fast_field_readers = U32FastFieldsReader::open(source).unwrap();
|
||||
|
||||
+11
-18
@@ -12,6 +12,7 @@ use directory::{WritePtr, RAMDirectory, Directory};
|
||||
use fastfield::FastFieldSerializer;
|
||||
use fastfield::U32FastFieldsWriter;
|
||||
use common::bitpacker::compute_num_bits;
|
||||
use common::bitpacker::BitUnpacker;
|
||||
|
||||
|
||||
lazy_static! {
|
||||
@@ -23,11 +24,9 @@ lazy_static! {
|
||||
|
||||
pub struct U32FastFieldReader {
|
||||
_data: ReadOnlySource,
|
||||
data_ptr: *const u8,
|
||||
bit_unpacker: BitUnpacker,
|
||||
min_val: u32,
|
||||
max_val: u32,
|
||||
num_bits: u32,
|
||||
mask: u32,
|
||||
}
|
||||
|
||||
impl U32FastFieldReader {
|
||||
@@ -47,34 +46,28 @@ impl U32FastFieldReader {
|
||||
pub fn open(data: ReadOnlySource) -> io::Result<U32FastFieldReader> {
|
||||
let min_val;
|
||||
let amplitude;
|
||||
let max_val;
|
||||
{
|
||||
let mut cursor = data.as_slice();
|
||||
min_val = try!(u32::deserialize(&mut cursor));
|
||||
amplitude = try!(u32::deserialize(&mut cursor));
|
||||
max_val = min_val + amplitude;
|
||||
}
|
||||
let num_bits = compute_num_bits(amplitude);
|
||||
let mask = (1 << num_bits) - 1;
|
||||
let ptr: *const u8 = &(data.deref()[8 as usize]);
|
||||
let bit_unpacker = {
|
||||
let data_arr = &(data.deref()[8..]);
|
||||
BitUnpacker::new(data_arr, num_bits as usize)
|
||||
};
|
||||
Ok(U32FastFieldReader {
|
||||
_data: data,
|
||||
data_ptr: ptr,
|
||||
bit_unpacker: bit_unpacker,
|
||||
min_val: min_val,
|
||||
max_val: min_val + amplitude,
|
||||
num_bits: num_bits as u32,
|
||||
mask: mask,
|
||||
max_val: max_val,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn get(&self, doc: DocId) -> u32 {
|
||||
if self.num_bits == 0u32 {
|
||||
return self.min_val;
|
||||
}
|
||||
let addr = (doc * self.num_bits) / 8;
|
||||
let bit_shift = (doc * self.num_bits) - addr * 8; //doc - long_addr * self.num_in_pack;
|
||||
let val_unshifted_unmasked: u64 = unsafe { * (self.data_ptr.offset(addr as isize) as *const u64) };
|
||||
let val_shifted = (val_unshifted_unmasked >> bit_shift) as u32;
|
||||
self.min_val + (val_shifted & self.mask)
|
||||
|
||||
self.min_val + self.bit_unpacker.get(doc as usize)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+15
-42
@@ -1,9 +1,8 @@
|
||||
use common::BinarySerializable;
|
||||
use directory::WritePtr;
|
||||
use schema::Field;
|
||||
use common::bitpacker::compute_num_bits;
|
||||
use std::io;
|
||||
use std::io::{Write, Seek, SeekFrom};
|
||||
use common::bitpacker::{compute_num_bits, BitPacker};
|
||||
use std::io::{self, Write, Seek, SeekFrom};
|
||||
|
||||
|
||||
/// `FastFieldSerializer` is in charge of serializing
|
||||
@@ -30,16 +29,12 @@ pub struct FastFieldSerializer {
|
||||
write: WritePtr,
|
||||
written_size: usize,
|
||||
fields: Vec<(Field, u32)>,
|
||||
num_bits: u8,
|
||||
min_value: u32,
|
||||
field_open: bool,
|
||||
|
||||
mini_buffer_written: usize,
|
||||
mini_buffer: u32,
|
||||
bit_packer: BitPacker,
|
||||
}
|
||||
|
||||
|
||||
|
||||
impl FastFieldSerializer {
|
||||
/// Constructor
|
||||
pub fn new(mut write: WritePtr) -> io::Result<FastFieldSerializer> {
|
||||
@@ -49,12 +44,9 @@ impl FastFieldSerializer {
|
||||
write: write,
|
||||
written_size: written_size,
|
||||
fields: Vec::new(),
|
||||
num_bits: 0u8,
|
||||
min_value: 0,
|
||||
field_open: false,
|
||||
|
||||
mini_buffer_written: 0,
|
||||
mini_buffer: 0u32,
|
||||
bit_packer: BitPacker::new(0),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -70,37 +62,23 @@ impl FastFieldSerializer {
|
||||
self.written_size += try!(min_value.serialize(write));
|
||||
let amplitude = max_value - min_value;
|
||||
self.written_size += try!(amplitude.serialize(write));
|
||||
self.num_bits = compute_num_bits(amplitude);
|
||||
if self.num_bits == 0 {
|
||||
// if num_bits == 0 we make sure that we still write one mini buffer
|
||||
// so that the reader code does not overflows and does not
|
||||
// contain a needless if statement.
|
||||
self.mini_buffer_written += 1;
|
||||
}
|
||||
let num_bits = compute_num_bits(amplitude);
|
||||
self.bit_packer = BitPacker::new(num_bits as usize);
|
||||
// TODO inspect whether it is required
|
||||
// if num_bits == 0 {
|
||||
// // if num_bits == 0 we make sure that we still write one mini buffer
|
||||
// // so that the reader code does not overflows and does not
|
||||
// // contain a needless if statement.
|
||||
// self.mini_buffer_written += 1;
|
||||
// }
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
/// Pushes a new value to the currently open u32 fast field.
|
||||
pub fn add_val(&mut self, val: u32) -> io::Result<()> {
|
||||
let write: &mut Write = &mut self.write;
|
||||
let val_to_write: u32 = val - self.min_value;
|
||||
if self.mini_buffer_written + self.num_bits as usize > 32 {
|
||||
self.mini_buffer |= val_to_write.wrapping_shl(self.mini_buffer_written as u32);
|
||||
self.written_size += try!(self.mini_buffer.serialize(write));
|
||||
// overflow of the shift operand is guarded here by the if case.
|
||||
self.mini_buffer = val_to_write.wrapping_shr(32u32 - self.mini_buffer_written as u32);
|
||||
self.mini_buffer_written = self.mini_buffer_written + (self.num_bits as usize) - 32 ;
|
||||
}
|
||||
else {
|
||||
self.mini_buffer |= val_to_write << self.mini_buffer_written;
|
||||
self.mini_buffer_written += self.num_bits as usize;
|
||||
if self.mini_buffer_written == 32 {
|
||||
self.written_size += try!(self.mini_buffer.serialize(write));
|
||||
self.mini_buffer_written = 0;
|
||||
self.mini_buffer = 0u32;
|
||||
}
|
||||
}
|
||||
self.bit_packer.write(val_to_write, &mut self.write)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -110,15 +88,10 @@ impl FastFieldSerializer {
|
||||
return Err(io::Error::new(io::ErrorKind::Other, "Current field is already closed"));
|
||||
}
|
||||
self.field_open = false;
|
||||
if self.mini_buffer_written > 0 {
|
||||
self.mini_buffer_written = 0;
|
||||
self.written_size += try!(self.mini_buffer.serialize(&mut self.write));
|
||||
}
|
||||
// adding some padding to make sure we
|
||||
// can read the last elements with our u64
|
||||
// cursor
|
||||
self.written_size += try!(0u32.serialize(&mut self.write));
|
||||
self.mini_buffer = 0;
|
||||
self.written_size += self.bit_packer.close(&mut self.write)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user