diff --git a/parquet/src/arrow/arrow_writer/levels.rs b/parquet/src/arrow/arrow_writer/levels.rs index 8374e905e1a0..d23873278ed3 100644 --- a/parquet/src/arrow/arrow_writer/levels.rs +++ b/parquet/src/arrow/arrow_writer/levels.rs @@ -336,51 +336,81 @@ impl LevelInfoBuilder { }) }; - let write_empty_slice = |child: &mut LevelInfoBuilder| { - child.visit_leaves(|leaf| { - let rep_levels = leaf.rep_levels.as_mut().unwrap(); - rep_levels.push(ctx.rep_level - 1); - let def_levels = leaf.def_levels.as_mut().unwrap(); - def_levels.push(ctx.def_level - 1); - }) + let write_null_run = |child: &mut LevelInfoBuilder, count: usize| { + if count > 0 { + child.visit_leaves(|leaf| { + leaf.rep_levels + .as_mut() + .unwrap() + .extend(std::iter::repeat_n(ctx.rep_level - 1, count)); + leaf.def_levels + .as_mut() + .unwrap() + .extend(std::iter::repeat_n(ctx.def_level - 2, count)); + }); + } }; - let write_null_slice = |child: &mut LevelInfoBuilder| { - child.visit_leaves(|leaf| { - let rep_levels = leaf.rep_levels.as_mut().unwrap(); - rep_levels.push(ctx.rep_level - 1); - let def_levels = leaf.def_levels.as_mut().unwrap(); - def_levels.push(ctx.def_level - 2); - }) + let write_empty_run = |child: &mut LevelInfoBuilder, count: usize| { + if count > 0 { + child.visit_leaves(|leaf| { + leaf.rep_levels + .as_mut() + .unwrap() + .extend(std::iter::repeat_n(ctx.rep_level - 1, count)); + leaf.def_levels + .as_mut() + .unwrap() + .extend(std::iter::repeat_n(ctx.def_level - 1, count)); + }); + } }; match nulls { Some(nulls) => { let null_offset = range.start; + let mut pending_nulls: usize = 0; + let mut pending_empties: usize = 0; + // TODO: Faster bitmask iteration (#1757) for (idx, w) in offsets.windows(2).enumerate() { let is_valid = nulls.is_valid(idx + null_offset); let start_idx = w[0].as_usize(); let end_idx = w[1].as_usize(); + if !is_valid { - write_null_slice(child) + write_empty_run(child, pending_empties); + pending_empties = 0; + pending_nulls += 1; } else if start_idx == end_idx { - write_empty_slice(child) + write_null_run(child, pending_nulls); + pending_nulls = 0; + pending_empties += 1; } else { - write_non_null_slice(child, start_idx, end_idx) + write_null_run(child, pending_nulls); + pending_nulls = 0; + write_empty_run(child, pending_empties); + pending_empties = 0; + write_non_null_slice(child, start_idx, end_idx); } } + write_null_run(child, pending_nulls); + write_empty_run(child, pending_empties); } None => { + let mut pending_empties: usize = 0; for w in offsets.windows(2) { let start_idx = w[0].as_usize(); let end_idx = w[1].as_usize(); if start_idx == end_idx { - write_empty_slice(child) + pending_empties += 1; } else { - write_non_null_slice(child, start_idx, end_idx) + write_empty_run(child, pending_empties); + pending_empties = 0; + write_non_null_slice(child, start_idx, end_idx); } } + write_empty_run(child, pending_empties); } } }