1use crate::row_index::RowSource;
17use polars::prelude::*;
18use std::collections::BTreeMap;
19use std::path::Path;
20use std::sync::Arc;
21
22const DAY_NS: i64 = 86_400_000_000_000;
24
25pub enum Bytes {
27 Mapped(memmap2::Mmap, std::fs::File),
29 Owned(Vec<u8>),
31}
32
33impl Bytes {
34 pub fn map(path: &Path) -> std::io::Result<Self> {
36 let file = std::fs::File::open(path)?;
37 if file.metadata()?.len() == 0 {
39 return Ok(Self::Owned(Vec::new()));
40 }
41 let map = unsafe { memmap2::Mmap::map(&file)? };
47 Ok(Self::Mapped(map, file))
48 }
49
50 pub fn as_slice(&self) -> &[u8] {
51 match self {
52 Self::Mapped(map, _) => map,
53 Self::Owned(bytes) => bytes,
54 }
55 }
56
57 pub fn still_whole(&self) -> PolarsResult<()> {
60 if let Self::Mapped(map, file) = self {
61 let len = file.metadata()?.len();
62 polars_ensure!(
63 len >= map.len() as u64,
64 ComputeError: "the file is now {len} bytes, shorter than the {} it had when it was opened; open it again",
65 map.len()
66 );
67 }
68 Ok(())
69 }
70
71 pub fn len(&self) -> usize {
72 self.as_slice().len()
73 }
74
75 pub fn is_empty(&self) -> bool {
76 self.len() == 0
77 }
78}
79
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
82pub enum Physical {
83 Unsigned(u8),
85 Signed(u8),
87 Float(u8),
89 BFloat16,
91 Bool,
93 Text,
95 Latin1,
97 Utf16 { big_endian: bool },
99 Utf32 { big_endian: bool },
102 Raw,
104}
105
106impl Physical {
107 pub fn width(self) -> Option<usize> {
109 match self {
110 Self::Unsigned(n) | Self::Signed(n) | Self::Float(n) => Some(n as usize),
111 Self::Bool => Some(1),
112 Self::BFloat16 => Some(2),
113 Self::Text | Self::Latin1 | Self::Utf16 { .. } | Self::Utf32 { .. } | Self::Raw => None,
114 }
115 }
116
117 pub fn is_integer(self) -> bool {
118 matches!(self, Self::Unsigned(_) | Self::Signed(_))
119 }
120
121 fn plain_dtype(self) -> DataType {
122 match self {
123 Self::Unsigned(1) => DataType::UInt8,
124 Self::Unsigned(2) => DataType::UInt16,
125 Self::Unsigned(3 | 4) => DataType::UInt32,
126 Self::Unsigned(_) => DataType::UInt64,
127 Self::Signed(1) => DataType::Int8,
128 Self::Signed(2) => DataType::Int16,
129 Self::Signed(3 | 4) => DataType::Int32,
130 Self::Signed(_) => DataType::Int64,
131 Self::Float(2 | 4) | Self::BFloat16 => DataType::Float32,
132 Self::Float(_) => DataType::Float64,
133 Self::Bool => DataType::Boolean,
134 Self::Text | Self::Latin1 | Self::Utf16 { .. } | Self::Utf32 { .. } => DataType::String,
135 Self::Raw => DataType::Binary,
136 }
137 }
138}
139
140#[derive(Debug, Clone, Copy, PartialEq)]
142pub enum Null {
143 Min,
145 Max,
147 NaN,
149 Value(i128),
151}
152
153#[derive(Debug, Clone, PartialEq)]
155pub enum Logical {
156 Plain,
158 Timestamp {
161 unit: TimeUnit,
162 multiplier: i64,
163 epoch: i64,
164 },
165 FloatTimestamp { ns_per_unit: f64, epoch_ns: i64 },
167 Duration { unit: TimeUnit, multiplier: i64 },
169 Days { epoch_days: i32 },
171 Yyyymmdd,
173 TimeOfDay {
176 ns_per_unit: i64,
177 date_ns: Option<i64>,
178 },
179 Decimal { scale: usize },
181 Linear { factor: f64, offset: f64 },
183 Enum(Arc<BTreeMap<i64, String>>),
185 Lookup(Arc<Vec<String>>),
187}
188
189#[derive(Debug, Clone)]
191pub struct ColumnLayout {
192 pub name: PlSmallStr,
193 pub source: usize,
195 pub start: usize,
197 pub stride: usize,
199 pub width: usize,
201 pub count: usize,
203 pub physical: Physical,
204 pub big_endian: bool,
205 pub null: Option<Null>,
206 pub logical: Logical,
207}
208
209impl ColumnLayout {
210 pub fn new(name: &str, start: usize, stride: usize, physical: Physical, width: usize) -> Self {
212 Self {
213 name: name.into(),
214 source: 0,
215 start,
216 stride,
217 width,
218 count: 1,
219 physical,
220 big_endian: false,
221 null: None,
222 logical: Logical::Plain,
223 }
224 }
225
226 pub fn value_dtype(&self) -> DataType {
228 match &self.logical {
229 Logical::Timestamp { unit, .. } => DataType::Datetime(*unit, None),
230 Logical::FloatTimestamp { .. } => DataType::Datetime(TimeUnit::Nanoseconds, None),
231 Logical::Duration { unit, .. } => DataType::Duration(*unit),
232 Logical::Days { .. } | Logical::Yyyymmdd => DataType::Date,
233 Logical::TimeOfDay { date_ns: None, .. } => DataType::Time,
234 Logical::TimeOfDay {
235 date_ns: Some(_), ..
236 } => DataType::Datetime(TimeUnit::Nanoseconds, None),
237 Logical::Decimal { scale } => DataType::Decimal(38, *scale),
238 Logical::Linear { .. } => DataType::Float64,
239 Logical::Enum(_) => DataType::String,
240 Logical::Lookup(_) => DataType::from_categories(Categories::global()),
241 Logical::Plain => self.physical.plain_dtype(),
242 }
243 }
244
245 pub fn dtype(&self) -> DataType {
247 if self.count > 1 {
248 DataType::Array(Box::new(self.value_dtype()), self.count)
249 } else {
250 self.value_dtype()
251 }
252 }
253
254 fn cell_width(&self) -> Option<usize> {
256 self.width.checked_mul(self.count.max(1))
257 }
258
259 fn rows_in(&self, len: usize) -> usize {
261 let Some(cell) = self.cell_width() else {
262 return 0;
263 };
264 match len
265 .checked_sub(self.start)
266 .and_then(|room| room.checked_sub(cell))
267 {
268 Some(after_first) if self.stride > 0 => after_first / self.stride + 1,
269 _ => 0,
270 }
271 }
272
273 pub fn validate(&self) -> PolarsResult<()> {
275 polars_ensure!(
276 self.stride > 0 && self.width > 0 && self.count > 0,
277 ComputeError: "column {} has no width", self.name
278 );
279 polars_ensure!(
280 self.cell_width().is_some(),
281 ComputeError: "column {}: {} values of {} bytes is too many", self.name, self.count, self.width
282 );
283 if let Some(width) = self.physical.width() {
284 polars_ensure!(
285 width == self.width,
286 ComputeError: "column {} is {} bytes wide, not {width}", self.name, self.width
287 );
288 }
289 if let Physical::Unsigned(n) | Physical::Signed(n) = self.physical {
290 polars_ensure!(
291 (1..=8).contains(&n),
292 ComputeError: "column {}: an integer is 1 to 8 bytes", self.name
293 );
294 }
295 if let Physical::Float(n) = self.physical {
296 polars_ensure!(
297 n == 2 || n == 4 || n == 8,
298 ComputeError: "column {}: a float is 2, 4 or 8 bytes", self.name
299 );
300 }
301 Ok(())
302 }
303}
304
305pub fn read_unsigned(bytes: &[u8], big_endian: bool) -> u64 {
307 let mut value = 0u64;
308 if big_endian {
309 for b in bytes {
310 value = (value << 8) | u64::from(*b);
311 }
312 } else {
313 for b in bytes.iter().rev() {
314 value = (value << 8) | u64::from(*b);
315 }
316 }
317 value
318}
319
320pub fn read_signed(bytes: &[u8], big_endian: bool) -> i64 {
322 let raw = read_unsigned(bytes, big_endian);
323 let bits = bytes.len() * 8;
324 if bits == 0 || bits >= 64 {
325 raw as i64
326 } else {
327 ((raw << (64 - bits)) as i64) >> (64 - bits)
329 }
330}
331
332fn null_integer(null: Null, physical: Physical) -> Option<i128> {
334 let bits = u32::from(match physical {
335 Physical::Unsigned(n) | Physical::Signed(n) => n,
336 Physical::Bool => 1,
337 _ => return None,
338 }) * 8;
339 let signed = matches!(physical, Physical::Signed(_));
340 match null {
341 Null::Min if signed => Some(-(1i128 << (bits - 1))),
342 Null::Min => Some(0),
343 Null::Max if signed => Some((1i128 << (bits - 1)) - 1),
344 Null::Max => Some((1i128 << bits) - 1),
345 Null::Value(v) => Some(v),
346 Null::NaN => None,
347 }
348}
349
350pub fn decode(bytes: &[u8], column: &ColumnLayout, rows: usize) -> PolarsResult<Column> {
354 column.validate()?;
355 let fits = column.rows_in(bytes.len());
356 polars_ensure!(
357 rows <= fits,
358 ComputeError: "column {}: {rows} rows asked for, {fits} in {} bytes", column.name, bytes.len()
359 );
360 decode_strided(bytes, column, rows, |row| row)
361}
362
363pub fn decode_rows(bytes: &[u8], column: &ColumnLayout, index: &IdxCa) -> PolarsResult<Column> {
366 column.validate()?;
367 let rows = crate::row_index::checked(index, column.rows_in(bytes.len()))?;
368 decode_strided(bytes, column, rows.len(), |i| rows[i] as usize)
369}
370
371pub fn decode_at(bytes: &[u8], column: &ColumnLayout, records: &[usize]) -> PolarsResult<Column> {
376 ColumnLayout {
378 stride: column.stride.max(1),
379 ..column.clone()
380 }
381 .validate()?;
382 let cell = column.cell_width().unwrap_or(usize::MAX);
383 for &record in records {
384 let end = record
385 .checked_add(column.start)
386 .and_then(|start| start.checked_add(cell));
387 polars_ensure!(
388 end.is_some_and(|end| end <= bytes.len()),
389 OutOfBounds: "column {}: a record at {record} runs past the {} bytes on hand", column.name, bytes.len()
390 );
391 }
392 let start = column.start;
393 decode_cells(bytes, column, records.len(), |i| records[i] + start)
394}
395
396fn decode_strided(
399 bytes: &[u8],
400 column: &ColumnLayout,
401 rows: usize,
402 row: impl Fn(usize) -> usize,
403) -> PolarsResult<Column> {
404 let (start, stride) = (column.start, column.stride);
405 decode_cells(bytes, column, rows, move |i| start + row(i) * stride)
406}
407
408fn decode_cells(
411 bytes: &[u8],
412 column: &ColumnLayout,
413 rows: usize,
414 cell: impl Fn(usize) -> usize,
415) -> PolarsResult<Column> {
416 let count = column.count.max(1);
417 let values = rows
418 .checked_mul(count)
419 .ok_or_else(|| polars_err!(ComputeError: "column {}: too many values", column.name))?;
420 let at = move |i: usize| {
422 let start = cell(i / count) + (i % count) * column.width;
423 &bytes[start..start + column.width]
424 };
425 let name = column.name.clone();
426 let flat = decode_values(column, values, at)?.with_name(name.clone());
427 if count == 1 {
428 return Ok(flat.into_column());
429 }
430 flat.reshape_array(&[
431 ReshapeDimension::new(rows as i64),
432 ReshapeDimension::new(count as i64),
433 ])
434 .map(|s| s.with_name(name).into_column())
435}
436
437fn decode_values<'a>(
439 column: &ColumnLayout,
440 n: usize,
441 at: impl Fn(usize) -> &'a [u8],
442) -> PolarsResult<Series> {
443 let big = column.big_endian;
444 let name = PlSmallStr::EMPTY;
445 macro_rules! native {
446 ($ty:ty) => {{
447 const N: usize = std::mem::size_of::<$ty>();
448 (0..n)
449 .map(|i| {
450 let raw: [u8; N] = at(i).try_into().expect("width checked at build");
451 if big {
452 <$ty>::from_be_bytes(raw)
453 } else {
454 <$ty>::from_le_bytes(raw)
455 }
456 })
457 .collect::<Vec<$ty>>()
458 }};
459 }
460 if column.logical == Logical::Plain && column.null.is_none() {
462 let fast = match column.physical {
463 Physical::Unsigned(1) => Some(Series::new(name.clone(), native!(u8))),
464 Physical::Unsigned(2) => Some(Series::new(name.clone(), native!(u16))),
465 Physical::Unsigned(4) => Some(Series::new(name.clone(), native!(u32))),
466 Physical::Unsigned(8) => Some(Series::new(name.clone(), native!(u64))),
467 Physical::Signed(1) => Some(Series::new(name.clone(), native!(i8))),
468 Physical::Signed(2) => Some(Series::new(name.clone(), native!(i16))),
469 Physical::Signed(4) => Some(Series::new(name.clone(), native!(i32))),
470 Physical::Signed(8) => Some(Series::new(name.clone(), native!(i64))),
471 Physical::Float(4) => Some(Series::new(name.clone(), native!(f32))),
472 Physical::Float(8) => Some(Series::new(name.clone(), native!(f64))),
473 Physical::Unsigned(3) => Some(Series::new(
476 name.clone(),
477 (0..n)
478 .map(|i| read_unsigned(at(i), big) as u32)
479 .collect::<Vec<u32>>(),
480 )),
481 Physical::Unsigned(5..=7) => Some(Series::new(
482 name.clone(),
483 (0..n)
484 .map(|i| read_unsigned(at(i), big))
485 .collect::<Vec<u64>>(),
486 )),
487 Physical::Signed(3) => Some(Series::new(
488 name.clone(),
489 (0..n)
490 .map(|i| read_signed(at(i), big) as i32)
491 .collect::<Vec<i32>>(),
492 )),
493 Physical::Signed(5..=7) => Some(Series::new(
494 name.clone(),
495 (0..n)
496 .map(|i| read_signed(at(i), big))
497 .collect::<Vec<i64>>(),
498 )),
499 _ => None,
500 };
501 if let Some(series) = fast {
502 return Ok(series);
503 }
504 }
505 match column.physical {
506 Physical::Unsigned(_) | Physical::Signed(_) | Physical::Bool => {
507 let signed = matches!(column.physical, Physical::Signed(_));
508 let sentinel = column
509 .null
510 .and_then(|null| null_integer(null, column.physical));
511 let ints: Vec<Option<i128>> = (0..n)
512 .map(|i| {
513 let bytes = at(i);
514 let v = if signed {
515 i128::from(read_signed(bytes, big))
516 } else {
517 i128::from(read_unsigned(bytes, big))
518 };
519 (Some(v) != sentinel).then_some(v)
520 })
521 .collect();
522 integers(column, ints)
523 }
524 Physical::Float(_) | Physical::BFloat16 => {
525 let physical = column.physical;
526 let floats: Vec<Option<f64>> = (0..n)
527 .map(|i| {
528 let raw = read_unsigned(at(i), big);
529 let v = match physical {
530 Physical::Float(2) => f64::from(half::f16::from_bits(raw as u16)),
531 Physical::BFloat16 => f64::from(half::bf16::from_bits(raw as u16)),
532 Physical::Float(4) => f64::from(f32::from_bits(raw as u32)),
533 _ => f64::from_bits(raw),
534 };
535 let null = match column.null {
536 Some(Null::NaN) => v.is_nan(),
537 Some(Null::Value(sentinel)) => v == sentinel as f64,
538 _ => false,
539 };
540 (!null).then_some(v)
541 })
542 .collect();
543 let width = physical.width().unwrap_or(8) as u8;
544 floats_of(column, width, floats)
545 }
546 Physical::Text => {
547 let values: StringChunked = (0..n).map(|i| Some(text(at(i)))).collect();
548 Ok(values.into_series())
549 }
550 Physical::Latin1 => {
551 let values: StringChunked = (0..n).map(|i| Some(latin1(at(i)))).collect();
552 Ok(values.into_series())
553 }
554 Physical::Utf16 { big_endian } => {
555 let values: StringChunked = (0..n).map(|i| Some(utf16(at(i), big_endian))).collect();
556 Ok(values.into_series())
557 }
558 Physical::Utf32 { big_endian } => {
559 let values: StringChunked = (0..n).map(|i| Some(utf32(at(i), big_endian))).collect();
560 Ok(values.into_series())
561 }
562 Physical::Raw => {
563 let values: BinaryChunked = (0..n).map(|i| Some(at(i))).collect();
564 Ok(values.into_series())
565 }
566 }
567}
568
569pub fn integers(column: &ColumnLayout, ints: Vec<Option<i128>>) -> PolarsResult<Series> {
573 let as_i64 = |v: Option<i128>| v.and_then(|v| i64::try_from(v).ok());
574 Ok(match &column.logical {
575 Logical::Plain => match column.physical {
576 Physical::Bool => ints
577 .into_iter()
578 .map(|v| v.map(|v| v != 0))
579 .collect::<BooleanChunked>()
580 .into_series(),
581 physical => {
582 let wide: Int128Chunked = ints.into_iter().collect();
583 wide.into_series().strict_cast(&physical.plain_dtype())?
585 }
586 },
587 Logical::Timestamp {
588 unit,
589 multiplier,
590 epoch,
591 } => ints
592 .into_iter()
593 .map(|v| as_i64(v)?.checked_mul(*multiplier)?.checked_add(*epoch))
594 .collect::<Int64Chunked>()
595 .into_datetime(*unit, None)
596 .into_series(),
597 Logical::Duration { unit, multiplier } => ints
598 .into_iter()
599 .map(|v| as_i64(v)?.checked_mul(*multiplier))
600 .collect::<Int64Chunked>()
601 .into_duration(*unit)
602 .into_series(),
603 Logical::FloatTimestamp {
604 ns_per_unit,
605 epoch_ns,
606 } => float_timestamps(
607 ints.into_iter().map(|v| v.map(|v| v as f64)),
608 *ns_per_unit,
609 *epoch_ns,
610 ),
611 Logical::Days { epoch_days } => ints
612 .into_iter()
613 .map(|v| i32::try_from(as_i64(v)?).ok()?.checked_add(*epoch_days))
614 .collect::<Int32Chunked>()
615 .into_date()
616 .into_series(),
617 Logical::Yyyymmdd => ints
618 .into_iter()
619 .map(|v| yyyymmdd_days(as_i64(v)?))
620 .collect::<Int32Chunked>()
621 .into_date()
622 .into_series(),
623 Logical::TimeOfDay {
624 ns_per_unit,
625 date_ns,
626 } => {
627 let of_day = ints.into_iter().map(|v| {
628 let ns = as_i64(v)?.checked_mul(*ns_per_unit)?;
629 (0..DAY_NS).contains(&ns).then_some(ns)
630 });
631 match date_ns {
632 None => of_day.collect::<Int64Chunked>().into_time().into_series(),
633 Some(day) => of_day
634 .map(|ns| ns?.checked_add(*day))
635 .collect::<Int64Chunked>()
636 .into_datetime(TimeUnit::Nanoseconds, None)
637 .into_series(),
638 }
639 }
640 Logical::Decimal { scale } => ints
641 .into_iter()
642 .collect::<Int128Chunked>()
643 .into_decimal_unchecked(38, *scale)
644 .into_series(),
645 Logical::Linear { factor, offset } => ints
646 .into_iter()
647 .map(|v| v.map(|v| v as f64 * factor + offset))
648 .collect::<Float64Chunked>()
649 .into_series(),
650 Logical::Lookup(symbols) => ints
651 .into_iter()
652 .map(|v| {
653 let i = usize::try_from(v?).ok()?;
654 symbols.get(i).map(String::as_str)
655 })
656 .collect::<StringChunked>()
657 .into_series()
658 .cast(&DataType::from_categories(Categories::global()))?,
659 Logical::Enum(labels) => ints
660 .into_iter()
661 .map(|v| {
662 v.map(|code| {
663 i64::try_from(code)
664 .ok()
665 .and_then(|code| labels.get(&code).cloned())
666 .unwrap_or_else(|| code.to_string())
667 })
668 })
669 .collect::<StringChunked>()
670 .into_series(),
671 })
672}
673
674fn floats_of(column: &ColumnLayout, width: u8, floats: Vec<Option<f64>>) -> PolarsResult<Series> {
676 Ok(match &column.logical {
677 Logical::Linear { factor, offset } => floats
678 .into_iter()
679 .map(|v| v.map(|v| v * factor + offset))
680 .collect::<Float64Chunked>()
681 .into_series(),
682 Logical::FloatTimestamp {
683 ns_per_unit,
684 epoch_ns,
685 } => float_timestamps(floats.into_iter(), *ns_per_unit, *epoch_ns),
686 _ if width <= 4 => floats
687 .into_iter()
688 .map(|v| v.map(|v| v as f32))
689 .collect::<Float32Chunked>()
690 .into_series(),
691 _ => floats.into_iter().collect::<Float64Chunked>().into_series(),
692 })
693}
694
695fn float_timestamps(
696 values: impl Iterator<Item = Option<f64>>,
697 ns_per_unit: f64,
698 epoch_ns: i64,
699) -> Series {
700 values
701 .map(|v| {
702 let ns = (v? * ns_per_unit).round();
703 (ns.is_finite() && ns.abs() < 9.0e18).then(|| (ns as i64).checked_add(epoch_ns))?
705 })
706 .collect::<Int64Chunked>()
707 .into_datetime(TimeUnit::Nanoseconds, None)
708 .into_series()
709}
710
711fn yyyymmdd_days(v: i64) -> Option<i32> {
713 let (year, month, day) = (v / 10_000, (v / 100) % 100, v % 100);
714 let date = chrono::NaiveDate::from_ymd_opt(
715 i32::try_from(year).ok()?,
716 u32::try_from(month).ok()?,
717 u32::try_from(day).ok()?,
718 )?;
719 let epoch = chrono::NaiveDate::from_ymd_opt(1970, 1, 1)?;
720 i32::try_from((date - epoch).num_days()).ok()
721}
722
723pub fn text(bytes: &[u8]) -> String {
725 let end = bytes
726 .iter()
727 .rposition(|&b| b != 0 && b != b' ')
728 .map_or(0, |i| i + 1);
729 String::from_utf8_lossy(&bytes[..end]).into_owned()
730}
731
732pub fn latin1(bytes: &[u8]) -> String {
734 let end = bytes
735 .iter()
736 .rposition(|&b| b != 0 && b != b' ')
737 .map_or(0, |i| i + 1);
738 bytes[..end].iter().map(|&b| char::from(b)).collect()
739}
740
741pub fn utf16(bytes: &[u8], big_endian: bool) -> String {
744 let mut units: Vec<u16> = bytes
745 .as_chunks::<2>()
746 .0
747 .iter()
748 .map(|&pair| {
749 if big_endian {
750 u16::from_be_bytes(pair)
751 } else {
752 u16::from_le_bytes(pair)
753 }
754 })
755 .collect();
756 while units.last().is_some_and(|&u| u == 0 || u == 0x20) {
757 units.pop();
758 }
759 String::from_utf16_lossy(&units)
760}
761
762pub fn utf32(bytes: &[u8], big_endian: bool) -> String {
766 let (units, rest) = bytes.as_chunks::<4>();
767 let mut chars: Vec<char> = units
768 .iter()
769 .map(|&unit| {
770 let code = if big_endian {
771 u32::from_be_bytes(unit)
772 } else {
773 u32::from_le_bytes(unit)
774 };
775 char::from_u32(code).unwrap_or(char::REPLACEMENT_CHARACTER)
776 })
777 .collect();
778 while chars.last() == Some(&'\0') {
779 chars.pop();
780 }
781 let mut text: String = chars.into_iter().collect();
782 if !rest.is_empty() {
783 text.push(char::REPLACEMENT_CHARACTER);
784 }
785 text
786}
787
788pub fn hex(bytes: &[u8]) -> String {
790 const DIGITS: &[u8; 16] = b"0123456789abcdef";
791 let mut out = String::with_capacity(bytes.len() * 3);
792 for (i, b) in bytes.iter().enumerate() {
793 if i > 0 {
794 out.push(' ');
795 }
796 out.push(DIGITS[(b >> 4) as usize] as char);
797 out.push(DIGITS[(b & 0xf) as usize] as char);
798 }
799 out
800}
801
802pub struct FixedRecords {
804 sources: Vec<Arc<Bytes>>,
805 columns: Vec<ColumnLayout>,
806 rows: usize,
807 schema: SchemaRef,
808}
809
810impl FixedRecords {
811 pub fn new(
818 sources: Vec<Arc<Bytes>>,
819 columns: Vec<ColumnLayout>,
820 rows: usize,
821 ) -> PolarsResult<Self> {
822 let mut rows = rows.min(crate::row_index::MAX_ROWS);
823 for column in &columns {
824 column.validate()?;
825 let source = sources
826 .get(column.source)
827 .ok_or_else(|| polars_err!(ComputeError: "column {} has no source", column.name))?;
828 rows = rows.min(column.rows_in(source.len()));
829 }
830 let schema: Schema = columns
831 .iter()
832 .map(|c| Field::new(c.name.clone(), c.dtype()))
833 .collect();
834 polars_ensure!(
836 schema.len() == columns.len(),
837 Duplicate: "two columns have the same name"
838 );
839 Ok(Self {
840 sources,
841 columns,
842 rows,
843 schema: Arc::new(schema),
844 })
845 }
846
847 pub fn rows(&self) -> usize {
848 self.rows
849 }
850
851 pub fn schema(&self) -> SchemaRef {
852 self.schema.clone()
853 }
854
855 pub fn columns(&self) -> &[ColumnLayout] {
856 &self.columns
857 }
858
859 pub fn sources(&self) -> &[Arc<Bytes>] {
861 &self.sources
862 }
863
864 pub fn lazy(self: &Arc<Self>) -> LazyFrame {
866 crate::row_index::lazy(self)
867 }
868
869 pub fn window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
872 let start = start.min(self.rows);
873 let len = len.min(self.rows - start);
874 let columns = self
875 .columns
876 .iter()
877 .map(|c| ColumnLayout {
878 start: c.start + start * c.stride,
879 ..c.clone()
880 })
881 .collect();
882 Self::new(self.sources.clone(), columns, len)?.collect(len)
883 }
884
885 pub fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
887 let rows = rows.min(self.rows);
888 let columns = self
889 .columns
890 .iter()
891 .map(|c| self.decode_column(c, rows))
892 .collect::<PolarsResult<Vec<_>>>()?;
893 DataFrame::new(rows, columns)
894 }
895
896 fn decode_column(&self, column: &ColumnLayout, rows: usize) -> PolarsResult<Column> {
897 let source = &self.sources[column.source];
898 source.still_whole()?;
899 decode(source.as_slice(), column, rows)
900 }
901}
902
903impl RowSource for FixedRecords {
904 fn height(&self) -> usize {
905 self.rows
906 }
907
908 fn schema(&self) -> SchemaRef {
909 self.schema.clone()
910 }
911
912 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
913 let column = &self.columns[column];
914 let source = &self.sources[column.source];
915 source.still_whole()?;
916 decode_rows(source.as_slice(), column, index)
917 }
918}
919
920impl crate::pushdown::Windowed for FixedRecords {
921 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
922 Ok(FixedRecords::window(self, start, len)?.lazy())
923 }
924}
925
926#[cfg(test)]
927mod tests {
928 use super::*;
929
930 fn records(bytes: Vec<u8>, columns: Vec<ColumnLayout>, rows: usize) -> Arc<FixedRecords> {
931 Arc::new(FixedRecords::new(vec![Arc::new(Bytes::Owned(bytes))], columns, rows).unwrap())
932 }
933
934 fn column(name: &str, start: usize, stride: usize, physical: Physical) -> ColumnLayout {
935 ColumnLayout::new(
936 name,
937 start,
938 stride,
939 physical,
940 physical.width().unwrap_or(stride),
941 )
942 }
943
944 #[test]
948 #[ignore = "a timing, not a check"]
949 fn time_integer_widths() {
950 const ROWS: usize = 10_000_000;
951 const STRIDE: usize = 16;
952 let bytes: Vec<u8> = (0..ROWS * STRIDE).map(|i| (i * 31 % 251) as u8).collect();
953 for physical in [
954 Physical::Unsigned(2),
955 Physical::Unsigned(3),
956 Physical::Unsigned(4),
957 Physical::Signed(3),
958 Physical::Unsigned(5),
959 Physical::Signed(6),
960 Physical::Unsigned(8),
961 ] {
962 for big_endian in [false, true] {
963 let layout = ColumnLayout {
964 big_endian,
965 ..column("v", 1, STRIDE, physical)
966 };
967 let started = std::time::Instant::now();
968 let decoded = decode(&bytes, &layout, ROWS).unwrap();
969 println!(
970 "{physical:?} {}: {:?} ({})",
971 if big_endian { "be" } else { "le" },
972 started.elapsed(),
973 decoded.dtype()
974 );
975 }
976 }
977 }
978
979 #[test]
980 fn strided_columns_decode_in_both_byte_orders() {
981 let mut bytes = Vec::new();
983 for (a, b) in [(1u16, -2i32), (3, -4)] {
984 bytes.extend(a.to_le_bytes());
985 bytes.extend(b.to_le_bytes());
986 }
987 let lf = records(
988 bytes.clone(),
989 vec![
990 column("a", 0, 6, Physical::Unsigned(2)),
991 column("b", 2, 6, Physical::Signed(4)),
992 ],
993 usize::MAX,
994 )
995 .lazy();
996 let df = lf.collect().unwrap();
997 assert_eq!(df.height(), 2);
998 assert_eq!(
999 df.column("b").unwrap().i32().unwrap().to_vec(),
1000 [Some(-2), Some(-4)]
1001 );
1002 let mut big = column("a", 0, 6, Physical::Unsigned(2));
1003 big.big_endian = true;
1004 let df = records(bytes, vec![big], usize::MAX).collect(9).unwrap();
1005 assert_eq!(
1006 df.column("a").unwrap().u16().unwrap().to_vec(),
1007 [Some(256), Some(768)]
1008 );
1009 }
1010
1011 #[test]
1012 fn odd_widths_sign_extend_and_sentinels_read_null() {
1013 let bytes = vec![0xfe, 0xff, 0xff, 0xff, 0xff, 0x7f, 0x00, 0x00, 0x80];
1015 let mut s3 = column("s", 0, 3, Physical::Signed(3));
1016 let df = records(bytes.clone(), vec![s3.clone()], usize::MAX)
1017 .collect(9)
1018 .unwrap();
1019 assert_eq!(
1020 df.column("s").unwrap().i32().unwrap().to_vec(),
1021 [Some(-2), Some(8_388_607), Some(-8_388_608)]
1022 );
1023 s3.null = Some(Null::Min);
1024 let df = records(bytes, vec![s3], usize::MAX).collect(9).unwrap();
1025 assert_eq!(df.column("s").unwrap().null_count(), 1);
1026 let u5 = ColumnLayout {
1028 big_endian: true,
1029 ..column("u5", 0, 11, Physical::Unsigned(5))
1030 };
1031 let s6 = column("s6", 5, 11, Physical::Signed(6));
1032 let mut bytes = vec![0x01, 0, 0, 0, 0x02];
1033 bytes.extend((-3i64).to_le_bytes()[..6].iter());
1034 let df = records(bytes, vec![u5, s6], 1).collect(1).unwrap();
1035 assert_eq!(
1036 df.column("u5").unwrap().u64().unwrap().get(0),
1037 Some(0x01_0000_0002)
1038 );
1039 assert_eq!(df.column("s6").unwrap().i64().unwrap().get(0), Some(-3));
1040 let mut u2 = column("u", 0, 2, Physical::Unsigned(2));
1041 u2.null = Some(Null::Max);
1042 let df = records(vec![0xff, 0xff, 1, 0], vec![u2], usize::MAX)
1043 .collect(9)
1044 .unwrap();
1045 assert_eq!(
1046 df.column("u").unwrap().u16().unwrap().to_vec(),
1047 [None, Some(1)]
1048 );
1049 }
1050
1051 #[test]
1052 fn a_count_makes_an_array_and_factor_offset_a_float() {
1053 let mut channels = column("ch", 0, 3, Physical::Unsigned(1));
1055 channels.count = 3;
1056 channels.logical = Logical::Linear {
1057 factor: 0.5,
1058 offset: -1.0,
1059 };
1060 let df = records(vec![0, 2, 4, 6, 8, 10], vec![channels], usize::MAX)
1061 .collect(9)
1062 .unwrap();
1063 let ch = df.column("ch").unwrap();
1064 assert_eq!(ch.dtype(), &DataType::Array(Box::new(DataType::Float64), 3));
1065 assert_eq!(ch.get(1).unwrap().to_string(), "[2.0, 3.0, 4.0]");
1066 }
1067
1068 #[test]
1069 fn dates_and_times_of_day() {
1070 let mut ymd = column("d", 0, 4, Physical::Unsigned(4));
1071 ymd.logical = Logical::Yyyymmdd;
1072 let mut bytes = 20240229u32.to_le_bytes().to_vec();
1073 bytes.extend(20241301u32.to_le_bytes());
1074 let df = records(bytes, vec![ymd], usize::MAX).collect(9).unwrap();
1075 let d = df.column("d").unwrap();
1076 assert_eq!(d.get(0).unwrap().to_string(), "2024-02-29");
1077 assert_eq!(d.null_count(), 1, "month 13 is no date");
1078 let mut tod = column("t", 0, 6, Physical::Unsigned(6));
1079 tod.big_endian = true;
1080 tod.logical = Logical::TimeOfDay {
1081 ns_per_unit: 1,
1082 date_ns: None,
1083 };
1084 let ns: u64 = 34_200_000_000_000; let df = records(ns.to_be_bytes()[2..].to_vec(), vec![tod], usize::MAX)
1086 .collect(9)
1087 .unwrap();
1088 assert_eq!(
1089 df.column("t").unwrap().get(0).unwrap().to_string(),
1090 "09:30:00"
1091 );
1092 }
1093
1094 #[test]
1095 fn projection_and_n_rows_reach_the_scan() {
1096 let bytes: Vec<u8> = (0u8..40).collect();
1097 let lf = records(
1098 bytes,
1099 vec![
1100 column("a", 0, 4, Physical::Unsigned(1)),
1101 column("b", 1, 4, Physical::Unsigned(1)),
1102 ],
1103 usize::MAX,
1104 )
1105 .lazy();
1106 let df = lf.clone().select([col("b")]).limit(3).collect().unwrap();
1107 assert_eq!(df.get_column_names(), ["b"]);
1108 assert_eq!(df.height(), 3);
1109 let count = lf.select([len()]).collect().unwrap();
1110 assert_eq!(
1111 count.column("len").unwrap().get(0).unwrap(),
1112 AnyValue::UInt32(10)
1113 );
1114 }
1115
1116 #[test]
1117 fn a_window_starts_where_it_is_asked() {
1118 let bytes: Vec<u8> = (0u8..40).collect();
1119 let records = records(
1120 bytes,
1121 vec![column("a", 0, 4, Physical::Unsigned(1))],
1122 usize::MAX,
1123 );
1124 let df = records.window(8, 5).unwrap();
1125 assert_eq!(
1126 df.column("a").unwrap().u8().unwrap().to_vec(),
1127 [Some(32), Some(36)]
1128 );
1129 assert_eq!(records.window(99, 5).unwrap().height(), 0);
1130 }
1131
1132 #[test]
1134 fn a_read_past_the_bytes_is_an_error_not_a_panic() {
1135 let u2 = column("u", 0, 2, Physical::Unsigned(2));
1136 assert!(decode(&[1, 0, 2, 0], &u2, 2).is_ok());
1137 assert!(decode(&[1, 0, 2, 0], &u2, 3).is_err());
1138 let mut huge = u2.clone();
1139 huge.count = usize::MAX;
1140 assert!(decode(&[0; 4], &huge, 0).is_err());
1141 assert!(
1142 FixedRecords::new(vec![Arc::new(Bytes::Owned(vec![0; 4]))], vec![huge], 1).is_err()
1143 );
1144 let mut far = u2;
1145 far.start = usize::MAX;
1146 let records = records(vec![0; 4], vec![far], usize::MAX);
1147 assert_eq!(records.rows(), 0);
1148 }
1149
1150 #[test]
1154 fn a_file_that_shrank_is_refused() {
1155 let dir = tempfile::tempdir().unwrap();
1156 let path = dir.path().join("f.bin");
1157 std::fs::write(&path, [7u8; 64]).unwrap();
1158 let bytes = Arc::new(Bytes::map(&path).unwrap());
1159 let records = Arc::new(
1160 FixedRecords::new(
1161 vec![bytes],
1162 vec![column("a", 0, 1, Physical::Unsigned(1))],
1163 usize::MAX,
1164 )
1165 .unwrap(),
1166 );
1167 assert_eq!(records.collect(64).unwrap().height(), 64);
1168 let cut = std::fs::OpenOptions::new()
1169 .write(true)
1170 .open(&path)
1171 .unwrap()
1172 .set_len(8);
1173 if cfg!(windows) {
1174 assert_eq!(cut.unwrap_err().raw_os_error(), Some(1224));
1175 assert_eq!(records.collect(64).unwrap().height(), 64);
1176 return;
1177 }
1178 cut.unwrap();
1179 let err = records.lazy().collect().unwrap_err();
1180 assert!(err.to_string().contains("shorter"), "{err}");
1181 }
1182
1183 fn numbered(rows: u32) -> Arc<FixedRecords> {
1185 let mut bytes = Vec::new();
1186 for i in 0..rows {
1187 bytes.extend(i.to_le_bytes());
1188 bytes.extend(((i % 7) as i16 - 3).to_le_bytes());
1189 bytes.push((i % 2) as u8);
1190 }
1191 records(
1192 bytes,
1193 vec![
1194 column("id", 0, 7, Physical::Unsigned(4)),
1195 column("group", 4, 7, Physical::Signed(2)),
1196 column("flag", 6, 7, Physical::Bool),
1197 ],
1198 usize::MAX,
1199 )
1200 }
1201
1202 #[test]
1205 fn queries_stream_and_agree_with_the_in_memory_engine() {
1206 let lf = numbered(1_000).lazy();
1207 let queries = [
1208 lf.clone()
1209 .filter(col("flag").and(col("group").gt(lit(0i16))))
1210 .select([len()]),
1211 lf.clone()
1212 .sort(
1213 ["group", "id"],
1214 SortMultipleOptions::default().with_order_descending(true),
1215 )
1216 .limit(5),
1217 lf.clone()
1218 .group_by([col("group")])
1219 .agg([col("id").sum(), len()])
1220 .sort(["group"], Default::default()),
1221 lf.clone().slice(990, 50),
1222 ];
1223 for query in queries {
1224 let streamed = crate::statistics::collect_lazy(query.clone(), true).unwrap();
1225 let in_memory = query.collect().unwrap();
1226 assert!(
1227 streamed.equals_missing(&in_memory),
1228 "{streamed}\n{in_memory}"
1229 );
1230 }
1231 let groups = crate::statistics::collect_lazy(
1232 lf.group_by([col("group")])
1233 .agg([len()])
1234 .sort(["group"], Default::default()),
1235 true,
1236 )
1237 .unwrap();
1238 assert_eq!(groups.height(), 7);
1239 assert_eq!(
1240 groups.column("group").unwrap().i16().unwrap().get(0),
1241 Some(-3)
1242 );
1243 }
1244
1245 #[test]
1249 fn a_deep_slice_reads_its_own_rows() {
1250 let records = numbered(100_000);
1251 let window = records.window(99_990, 50).unwrap();
1252 assert_eq!(window.height(), 10);
1253 for streaming in [false, true] {
1254 let sliced =
1255 crate::statistics::collect_lazy(records.lazy().slice(99_990, 50), streaming)
1256 .unwrap();
1257 assert!(window.equals_missing(&sliced), "{sliced}");
1258 }
1259 assert_eq!(
1260 window.column("id").unwrap().u32().unwrap().get(0),
1261 Some(99_990)
1262 );
1263 let got = Arc::new(std::sync::Mutex::new(0usize));
1264 let seen = got.clone();
1265 let sink = records
1266 .lazy()
1267 .filter(col("flag"))
1268 .sink_batches(
1269 PlanCallback::new(move |batch: DataFrame| {
1270 *seen.lock().unwrap() += batch.height();
1271 Ok(false)
1272 }),
1273 true,
1274 None,
1275 )
1276 .unwrap();
1277 crate::statistics::collect_lazy(sink, true).unwrap();
1278 assert_eq!(*got.lock().unwrap(), 50_000);
1279 }
1280
1281 #[test]
1284 fn a_sort_by_a_decimal_column_reads_its_first_page() {
1285 let mut price = column("price", 0, 4, Physical::Unsigned(4));
1286 price.logical = Logical::Decimal { scale: 2 };
1287 let bytes: Vec<u8> = (0u32..1_000).flat_map(|v| v.to_le_bytes()).collect();
1288 let lf = records(bytes, vec![price], usize::MAX).lazy();
1289 let page = crate::statistics::collect_lazy(
1290 lf.sort(
1291 ["price"],
1292 SortMultipleOptions::default().with_order_descending(true),
1293 )
1294 .slice(0, 3),
1295 true,
1296 )
1297 .unwrap();
1298 assert_eq!(
1299 page.column("price").unwrap().get(0).unwrap().to_string(),
1300 "9.99"
1301 );
1302 }
1303
1304 #[test]
1306 fn two_columns_of_one_name_are_refused() {
1307 let bytes = Arc::new(Bytes::Owned(vec![0; 8]));
1308 let a = column("a", 0, 2, Physical::Unsigned(1));
1309 let Err(err) = FixedRecords::new(vec![bytes], vec![a.clone(), a], usize::MAX) else {
1310 panic!("two columns named a were taken");
1311 };
1312 assert!(err.to_string().contains("same name"), "{err}");
1313 }
1314
1315 #[test]
1317 fn decoding_rows_checks_the_index() {
1318 let u1 = column("u", 0, 1, Physical::Unsigned(1));
1319 let index = IdxCa::from_slice("i".into(), &[3, 0, 3]);
1320 let col = decode_rows(&[5, 6, 7, 8], &u1, &index).unwrap();
1321 assert_eq!(col.u8().unwrap().to_vec(), [Some(8), Some(5), Some(8)]);
1322 assert!(decode_rows(&[5, 6, 7], &u1, &index).is_err());
1323 }
1324}