use std::cell::Cell;
use std::collections::HashSet;
use std::fs::File;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use crossbeam_channel::{Receiver, Sender};
use arrow_array::{Array, RecordBatch, UInt32Array};
use arrow_schema::{DataType, Field, Schema, SchemaRef};
use arrow_select::take::take;
use geo::{Area, Geometry};
use geoarrow::array::from_arrow_array;
use rayon::prelude::*;
use crate::batch_processor::{extract_geometries_from_array, extract_geometries_opt_from_array};
use crate::input_set::{ConvertSource, ReadPlan, RowGroupSelection};
use super::accumulate::{is_carrier, level_accumulates, tiny_polygon_carriers, AccumulateLevel};
use super::assign::{apply_density_budget, assign_levels_bounded, AssignFeature, FeatureKind};
use super::cluster::{ClusterEntry, ClusterTables};
use super::coalesce::CoalesceInput;
use super::convert::{
append_coalesced_count_field, append_point_count_field, apply_cluster_columns,
apply_coalesced_count, build_generalization, build_level_batch, build_level_coalesce_table,
build_source_schema, class_ranking_provenance, coalesce_effective, coalesce_level_chains,
count_vertices, encode_concurrency_for, extract_class_ranks, extract_sort_keys,
fill_level_bytes, find_geometry_column, mixed_geometry_field, overture_road_ranking,
record_level_outcome, resolve_reserved_column_collisions, scan_feature,
validate_cluster_schema, validate_coalesce_schema, warn_plan_skipped_levels, ClassRanking,
CoalesceTable, ConvertError, ConvertOptions, ConvertReport, GroupInterner, SkippedLevelReport,
KNOWN_ROAD_CLASSES, ROAD_VOCAB_MIN_DISTINCT,
};
use super::level::{Crs, Mode, RankingProvenance};
use super::pipe::scoped_pipe;
use super::pipeline;
use super::simplify::{
carrier_square, full_resolution_fallback_count, simplify_cascade, simplify_step,
validation_skip_count, CascadeStep, CollapseMode, Representation, Simplified, SimplifyOptions,
};
use super::writer::{LevelSpec, LevelWriteOutcome, OverviewWriter, OverviewWriterOptions};
const UNASSIGNED_LEVEL: u8 = u8::MAX;
struct EmitLevel {
orig: u8,
gsd: f64,
zoom: Option<u8>,
hint: usize,
}
#[derive(Clone, Copy)]
pub(crate) enum Pass2Strategy {
#[cfg_attr(not(test), allow(dead_code))]
Serial,
Pipelined,
}
fn log_validation_skips(skips_before: u64) {
let skips = validation_skip_count() - skips_before;
if skips > 0 {
log::info!(
"[convert] {skips} oversized RDP candidate(s) skipped exact \
validity checking and were assumed valid (#242; geometry \
validity is not an overviews conformance requirement)"
);
}
}
fn stage_input_pass0(
source: &ConvertSource,
selected_row_groups: Option<&RowGroupSelection>,
row_groups_read: usize,
) {
if !source.is_remote() {
return;
}
let t_stage = Instant::now();
match source.stage_selected(selected_row_groups) {
Ok(()) => log::info!(
"[convert] staged {row_groups_read} selected row group(s) to local \
disk in {:.1}s",
t_stage.elapsed().as_secs_f64()
),
Err(e) => log::warn!(
"[convert] input staging failed ({e}); passes will read over the \
network (uncoalesced)"
),
}
}
fn resolve_and_log_in_flight_batches(requested: usize) -> usize {
let in_flight = super::convert::resolve_in_flight_batches(requested);
log::info!(
"[convert] pass 2 parallelism: {in_flight} read batch(es) in flight ({} core(s) detected)",
std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(0)
);
in_flight
}
fn current_rss_mib() -> Option<f64> {
memory_stats::memory_stats().map(|m| m.physical_mem as f64 / (1024.0 * 1024.0))
}
fn log_phase_rss(phase: &str, peak_mib: &mut f64) {
if let Some(rss) = current_rss_mib() {
if rss > *peak_mib {
*peak_mib = rss;
}
log::info!("[rss] {phase}: {rss:.0} MiB");
}
}
fn partition_emitted_levels(
level_specs: &[(f64, Option<u8>)],
counts: &[usize],
) -> (Vec<EmitLevel>, Vec<SkippedLevelReport>) {
let skipped = level_specs
.iter()
.enumerate()
.filter(|&(l, _)| counts[l] == 0)
.map(|(l, &(gsd, zoom))| SkippedLevelReport {
planned_level: l,
gsd,
zoom,
})
.collect();
let emitted = level_specs
.iter()
.enumerate()
.filter(|&(l, _)| counts[l] > 0)
.map(|(l, &(gsd, zoom))| EmitLevel {
orig: l as u8,
gsd,
zoom,
hint: counts[l],
})
.collect();
(emitted, skipped)
}
pub(super) fn build_level_schemas(
input_schema: &Schema,
geom_idx: usize,
geom_name: &str,
options: &ConvertOptions,
) -> (Schema, Schema, Schema) {
let geom_out_field = mixed_geometry_field(geom_name);
let source_schema = build_source_schema(input_schema, geom_idx, geom_out_field);
let cluster_schema = if options.cluster {
append_point_count_field(&source_schema)
} else {
source_schema.clone()
};
let out_schema = if options.coalesce_lines {
append_coalesced_count_field(&cluster_schema)
} else {
cluster_schema.clone()
};
(source_schema, cluster_schema, out_schema)
}
pub(super) fn build_writer_options(
writer_levels: Vec<LevelSpec>,
emitted_gsds: &[f64],
crs: Crs,
ranking_provenance: RankingProvenance,
renames: &[(String, String)],
options: &ConvertOptions,
) -> OverviewWriterOptions {
let mut writer_opts = OverviewWriterOptions::new(options.mode, writer_levels);
writer_opts.max_row_group_size = options.max_row_group_size;
writer_opts.row_group_size_policy = options.row_group_size_policy;
writer_opts.full_column_stats = options.full_column_stats;
writer_opts.cogp_compat_key = options.cogp_compat_key;
writer_opts.encode_concurrency = encode_concurrency_for(options.profile);
writer_opts.generalization = Some(build_generalization(
emitted_gsds,
crs,
options,
ranking_provenance,
renames,
));
writer_opts
}
fn select_row_groups_streaming(
source: &ConvertSource,
bbox_units: Option<&[f64; 4]>,
filter: Option<&super::filter::BoundFilter>,
) -> Result<Option<RowGroupSelection>, ConvertError> {
let bbox_selection: Option<RowGroupSelection> = match bbox_units {
Some(bb) => Some(source.select_row_groups(bb)?),
None => None,
};
let filter_selection: Option<RowGroupSelection> = match filter {
Some(f) => Some(source.select_row_groups_matching(f)?),
None => None,
};
Ok(match (bbox_selection, filter_selection) {
(Some(a), Some(b)) => Some(a.intersect(&b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
})
}
fn polygon_area_f32(g: &Geometry<f64>) -> f32 {
match g {
Geometry::Polygon(p) => p.unsigned_area() as f32,
Geometry::MultiPolygon(mp) => mp.unsigned_area() as f32,
_ => 0.0,
}
}
fn streaming_carriers(
options: &ConvertOptions,
features: &[AssignFeature],
feat_min_levels: &[u8],
areas: Vec<f32>,
level_gsds: &[f64],
level_reprs: &[Representation],
crs: Crs,
) -> Vec<Vec<usize>> {
let finest_planned = level_gsds.len().saturating_sub(1);
let acc_levels: Vec<AccumulateLevel> = level_gsds
.iter()
.enumerate()
.map(|(l, &gsd)| AccumulateLevel {
gsd_meters: gsd,
enabled: l != finest_planned
&& accumulator_enabled(options)
&& level_accumulates(options.simplify.collapse, level_reprs[l]),
})
.collect();
if !acc_levels.iter().any(|l| l.enabled) {
return vec![Vec::new(); level_gsds.len()];
}
let t = Instant::now();
let carriers = tiny_polygon_carriers(
features,
feat_min_levels,
&areas,
&acc_levels,
crs,
options.simplify.factor,
);
let total: usize = carriers.iter().map(Vec::len).sum();
log::info!(
"[convert] tiny-polygon accumulator: {total} placeholder square(s) across {} \
level(s) stand in for the polygons those levels dropped ({:.2}s)",
acc_levels.iter().filter(|l| l.enabled).count(),
t.elapsed().as_secs_f64()
);
carriers
}
fn accumulator_enabled(options: &ConvertOptions) -> bool {
matches!(options.mode, Mode::Duplicating)
&& (options.simplify.collapse == CollapseMode::Square
|| options
.representation
.iter()
.any(|b| b.repr == Representation::Square))
}
fn build_pass2_coalesce_tables(
coalesce_scratch: Option<&CoalesceScratch>,
emitted: &[EmitLevel],
finest: usize,
crs: Crs,
options: &ConvertOptions,
) -> Vec<Option<CoalesceTable>> {
match coalesce_scratch {
Some(scratch) => {
log::info!(
"[convert] building coalesce chain tables for {} level(s)",
emitted.len()
);
let inputs = scratch.inputs();
emitted
.par_iter()
.map(|e| {
let verbatim =
matches!(options.mode, Mode::Partitioning) || e.orig as usize == finest;
(!verbatim).then(|| {
build_level_coalesce_table(
&inputs,
e.orig as usize,
finest,
e.gsd,
crs,
options,
)
})
})
.collect()
}
None => std::iter::repeat_with(|| None)
.take(emitted.len())
.collect(),
}
}
fn build_cascade_chains(
emitted: &[EmitLevel],
finest: usize,
duplicating: bool,
options: &ConvertOptions,
) -> Vec<Vec<CascadeStep>> {
let repr_of =
|zoom: Option<u8>| super::convert::representation_for_zoom(&options.representation, zoom);
emitted
.iter()
.map(|e| {
let verbatim = matches!(options.mode, Mode::Partitioning) || e.orig as usize == finest;
if !duplicating || verbatim || !options.simplify.cascade {
return Vec::new();
}
let mut chain: Vec<CascadeStep> = emitted
.iter()
.filter(|f| (f.orig as usize) < finest && f.orig >= e.orig)
.map(|f| CascadeStep {
gsd_meters: f.gsd,
repr: repr_of(f.zoom),
})
.collect();
chain.reverse();
chain
})
.collect()
}
struct LevelCtxInputs<'a> {
source_schema: &'a Schema,
cluster_schema: &'a Schema,
out_schema: &'a Schema,
non_geom_cols: &'a [usize],
geom_idx: usize,
min_levels: &'a [u8],
acc_cols: &'a [usize],
kinds: Option<&'a [FeatureKind]>,
cluster_tables: Option<&'a ClusterTables>,
coalesce_tables: &'a [Option<CoalesceTable>],
cascade_chains: &'a [Vec<CascadeStep>],
carriers: &'a [Vec<usize>],
crs: Crs,
finest: usize,
duplicating: bool,
}
fn build_level_ctxs<'a>(
emitted: &[EmitLevel],
options: &'a ConvertOptions,
inputs: &LevelCtxInputs<'a>,
) -> Vec<LevelStreamCtx<'a>> {
let repr_of =
|zoom: Option<u8>| super::convert::representation_for_zoom(&options.representation, zoom);
emitted
.iter()
.enumerate()
.map(|(i, e)| {
let verbatim =
matches!(options.mode, Mode::Partitioning) || e.orig as usize == inputs.finest;
LevelStreamCtx {
source_schema: inputs.source_schema,
cluster_schema: inputs.cluster_schema,
out_schema: inputs.out_schema,
non_geom_cols: inputs.non_geom_cols,
geom_idx: inputs.geom_idx,
min_levels: inputs.min_levels,
orig_level: e.orig,
duplicating: inputs.duplicating,
verbatim,
gsd_m: e.gsd,
repr: repr_of(e.zoom),
crs: inputs.crs,
simplify: &options.simplify,
cluster_enabled: options.cluster,
cluster_table: inputs
.cluster_tables
.filter(|_| e.orig as usize != inputs.finest)
.map(|t| &t[e.orig as usize]),
acc_cols: inputs.acc_cols,
coalesce_enabled: options.coalesce_lines,
kinds: inputs.kinds,
coalesce_table: inputs.coalesce_tables[i].as_ref(),
cascade_chain: &inputs.cascade_chains[i],
carriers: &inputs.carriers[e.orig as usize],
}
})
.collect()
}
#[allow(clippy::too_many_arguments)]
fn run_pass2_levels(
writer: &mut OverviewWriter<File>,
ctxs: &[LevelStreamCtx<'_>],
hints: &[usize],
source: &ConvertSource,
options: &ConvertOptions,
selected_row_groups: Option<&RowGroupSelection>,
in_flight_batches: usize,
out_schema: &Schema,
num_rows: usize,
geom_bytes: u64,
strategy: Pass2Strategy,
) -> Result<Vec<(LevelWriteOutcome, usize, usize)>, ConvertError> {
let n = ctxs.len();
let level_stats: Vec<(LevelWriteOutcome, usize, usize)> = match strategy {
Pass2Strategy::Serial => ctxs
.iter()
.enumerate()
.map(|(i, ctx)| {
write_level_streaming(
writer,
i,
hints[i],
source,
options.read_batch_size,
in_flight_batches,
selected_row_groups,
ctx,
)
})
.collect::<Result<_, _>>()?,
Pass2Strategy::Pipelined => {
let buffered_rows: usize = hints[..n - 1].iter().sum();
let avg_geom_bytes = (num_rows > 0).then(|| geom_bytes / num_rows as u64);
let backing = pipeline::resolve_backing(
options.profile,
options.mode,
buffered_rows,
avg_geom_bytes,
);
log::info!(
"[convert] pass 2: building {n} overview level(s) from a \
single read (finest level streamed last)"
);
let mut stats = if n > 1 {
pipeline::run_pass2_buffered(
writer,
&ctxs[..n - 1],
&hints[..n - 1],
source,
options.read_batch_size,
selected_row_groups,
in_flight_batches,
backing,
out_schema,
)?
} else {
Vec::new()
};
stats.push(write_level_streaming(
writer,
n - 1,
hints[n - 1],
source,
options.read_batch_size,
in_flight_batches,
selected_row_groups,
&ctxs[n - 1],
)?);
stats
}
};
Ok(level_stats)
}
type ResolvedRanking = (
Option<Vec<Option<f64>>>,
RankingProvenance,
Option<Vec<u32>>,
);
fn resolve_ranking_tier(
plan: RankPlan,
explicit_keys: Vec<Option<f64>>,
confidence_keys: Vec<Option<f64>>,
explicit_groups: Vec<u32>,
collect_lines: bool,
feature_count: usize,
point_count: usize,
) -> ResolvedRanking {
let n = feature_count;
let size_fallback = || {
log::info!(
"overview ranking: no sort key specified or auto-detected; using size + \
deterministic-hash fallback"
);
RankingProvenance {
mode: "size-fallback".to_string(),
column: None,
ranks: None,
unknown_rank: None,
}
};
match plan {
RankPlan::ExplicitSort { name, .. } => {
log::info!("overview ranking: explicit numeric sort-key column {name:?}");
(
Some(explicit_keys),
RankingProvenance {
mode: "explicit-sort-key".to_string(),
column: Some(name),
ranks: None,
unknown_rank: None,
},
None,
)
}
RankPlan::ExplicitClass { ranking, .. } => {
log::info!(
"overview ranking: explicit class-ranking on column {:?} ({} named classes, unknown_rank={})",
ranking.column,
ranking.ranks.len(),
ranking.unknown_rank
);
(
Some(explicit_keys),
class_ranking_provenance("class-ranking", &ranking),
collect_lines.then_some(explicit_groups),
)
}
RankPlan::Auto { roads, confidence } => {
if let Some(cand) = roads
.into_iter()
.find(|c| c.found.len() >= ROAD_VOCAB_MIN_DISTINCT)
{
log::info!(
"overview ranking: auto-detected Overture road classes in column {:?}; \
applying built-in ranking (motorway > … > service > tail)",
cand.ranking.column
);
let prov = class_ranking_provenance("auto-overture-roads", &cand.ranking);
(Some(cand.keys), prov, collect_lines.then_some(cand.groups))
} else if let Some((_, col_name)) = confidence.filter(|_| n > 0 && point_count * 2 >= n)
{
log::info!(
"overview ranking: auto-detected Overture places confidence column {col_name:?} \
(numeric point ranking)"
);
(
Some(confidence_keys),
RankingProvenance {
mode: "auto-confidence".to_string(),
column: Some(col_name),
ranks: None,
unknown_rank: None,
},
None,
)
} else {
(None, size_fallback(), None)
}
}
RankPlan::SizeFallback => (None, size_fallback(), None),
}
}
struct LevelWriter {
writer: OverviewWriter<File>,
source_schema: Schema,
cluster_schema: Schema,
out_schema: Schema,
non_geom_cols: Vec<usize>,
}
#[allow(clippy::too_many_arguments)]
fn create_level_writer(
output_path: &Path,
input_schema: &Schema,
geom_idx: usize,
geom_field: &Field,
emitted: &[EmitLevel],
crs: Crs,
ranking_provenance: RankingProvenance,
renames: &[(String, String)],
options: &ConvertOptions,
) -> Result<LevelWriter, ConvertError> {
let geom_name = geom_field.name().clone();
let (source_schema, cluster_schema, out_schema) =
build_level_schemas(input_schema, geom_idx, &geom_name, options);
let writer_levels: Vec<LevelSpec> = emitted
.iter()
.map(|e| LevelSpec::new(e.gsd, e.zoom))
.collect();
let emitted_gsds: Vec<f64> = emitted.iter().map(|e| e.gsd).collect();
let writer_opts = build_writer_options(
writer_levels,
&emitted_gsds,
crs,
ranking_provenance,
renames,
options,
);
let writer = OverviewWriter::create(output_path, &out_schema, writer_opts)?;
let non_geom_cols: Vec<usize> = (0..input_schema.fields().len())
.filter(|&c| c != geom_idx)
.collect();
Ok(LevelWriter {
writer,
source_schema,
cluster_schema,
out_schema,
non_geom_cols,
})
}
struct WinnerTables {
level_specs: Vec<(f64, Option<u8>)>,
cluster_tables: Option<ClusterTables>,
kinds: Option<Vec<FeatureKind>>,
coalesce_scratch: Option<CoalesceScratch>,
min_levels: Vec<u8>,
counts: Vec<usize>,
carriers: Vec<Vec<usize>>,
finest: usize,
}
#[allow(clippy::too_many_arguments)]
fn resolve_winner_tables(
features: &mut Vec<AssignFeature>,
acc_values: Vec<Vec<Option<f64>>>,
areas: Vec<f32>,
coalesce_scratch: Option<CoalesceScratch>,
num_rows: usize,
crs: Crs,
options: &ConvertOptions,
peak_rss_mib: &mut f64,
) -> Result<WinnerTables, ConvertError> {
let level_specs = options.levels.resolve(options.gsd_base)?;
let level_gsds: Vec<f64> = level_specs.iter().map(|(g, _)| *g).collect();
let t_assign = Instant::now();
let level_reprs = super::convert::level_representations(&level_specs, &options.representation);
let assignment = assign_levels_bounded(
features,
&level_gsds,
&options.assign,
crs,
super::pipeline::pass1_grid_budget_bytes(options.profile),
&level_reprs,
);
let assign_secs = t_assign.elapsed().as_secs_f64();
let t_budget = Instant::now();
let assignment = if options.density.enabled {
apply_density_budget(
&assignment,
features,
&level_gsds,
&options.assign,
&options.density,
crs,
)
} else {
assignment
};
log::debug!(
"[profile] assignment+budget: {:.2}s (assign {:.2}s + budget {:.2}s)",
t_assign.elapsed().as_secs_f64(),
assign_secs,
t_budget.elapsed().as_secs_f64()
);
log::info!(
"[convert] level assignment complete: {} level(s) in {:.1}s",
level_gsds.len(),
t_assign.elapsed().as_secs_f64()
);
log_phase_rss("assignment+budget (winner tables)", peak_rss_mib);
let feat_min_levels: Vec<u8> = assignment.assignments.iter().map(|a| a.min_level).collect();
drop(assignment);
let carriers = streaming_carriers(
options,
features,
&feat_min_levels,
areas,
&level_gsds,
&level_reprs,
crs,
);
let cluster_tables: Option<ClusterTables> = if options.cluster {
let acc_feat: Vec<Vec<Option<f64>>> = acc_values
.iter()
.map(|vals| features.iter().map(|f| vals[f.index]).collect())
.collect();
Some(super::convert::build_verified_cluster_tables(
features,
&feat_min_levels,
&level_gsds,
&acc_feat,
crs,
options,
)?)
} else {
None
};
drop(acc_values);
let coalesce_on = coalesce_effective(
options,
coalesce_scratch.as_ref().map_or(0, |s| s.rows.len()),
);
let kinds: Option<Vec<FeatureKind>> = options.coalesce_lines.then(|| {
let mut k = vec![FeatureKind::Point; num_rows];
for f in features.iter() {
k[f.index] = f.kind;
}
k
});
let coalesce_scratch = coalesce_scratch.filter(|_| coalesce_on);
let num_levels = level_gsds.len();
let finest = num_levels.saturating_sub(1);
let mut min_levels = vec![UNASSIGNED_LEVEL; num_rows];
for (f, &ml) in features.iter().zip(&feat_min_levels) {
min_levels[f.index] = ml;
}
let mut hist = vec![0usize; num_levels];
for (f, &ml) in features.iter().zip(&feat_min_levels) {
if coalesce_scratch.is_some() && f.kind == FeatureKind::Line {
continue; }
hist[(ml as usize).min(finest)] += 1;
}
drop(feat_min_levels);
features.clear();
features.shrink_to_fit(); let mut counts: Vec<usize> = match options.mode {
Mode::Duplicating => hist
.iter()
.scan(0usize, |acc, &c| {
*acc += c;
Some(*acc)
})
.collect(),
Mode::Partitioning => hist,
};
for (count, level_carriers) in counts.iter_mut().zip(&carriers) {
*count += level_carriers.len();
}
if let Some(scratch) = &coalesce_scratch {
let inputs = scratch.inputs();
#[allow(clippy::needless_range_loop)]
for level in 0..num_levels {
if level == finest {
counts[level] += scratch.rows.len(); } else {
counts[level] +=
coalesce_level_chains(&inputs, level, finest, level_gsds[level], crs, options)
.len();
}
}
}
Ok(WinnerTables {
level_specs,
cluster_tables,
kinds,
coalesce_scratch,
min_levels,
counts,
carriers,
finest,
})
}
struct Preflight {
options: ConvertOptions,
input_schema: SchemaRef,
crs: Crs,
renames: Vec<(String, String)>,
geom_idx: usize,
geom_field: Field,
acc_cols: Vec<usize>,
bbox_units: Option<[f64; 4]>,
bound_filter: Option<super::filter::BoundFilter>,
selected_row_groups: Option<RowGroupSelection>,
row_groups_total: usize,
row_groups_read: usize,
}
fn convert_preflight(
source: &ConvertSource,
options: &ConvertOptions,
) -> Result<Preflight, ConvertError> {
let input_schema: SchemaRef = source.schema()?;
let kv = source.key_value_metadata()?;
let crs = super::convert::detect_crs_from_kv(kv.as_ref())?;
let mut resolved = options.clone();
let (input_schema, renames) = resolve_reserved_column_collisions(&input_schema, &mut resolved);
let options = &resolved;
let bound_filter = super::convert::bind_attribute_filter(options, &input_schema, &renames)?;
let row_groups_total = source.num_row_groups_total()?;
let bbox_units = options
.bbox
.map(|b| super::convert::bbox_to_crs_units(&b, crs));
let selected_row_groups =
select_row_groups_streaming(source, bbox_units.as_ref(), bound_filter.as_ref())?;
let row_groups_read = selected_row_groups
.as_ref()
.map_or(row_groups_total, RowGroupSelection::total_selected);
if selected_row_groups.is_some() {
let what = super::convert::pruning_label(options.bbox.is_some(), bound_filter.is_some());
log::info!("{what} filter: reading {row_groups_read}/{row_groups_total} input row groups");
}
super::convert::warn_full_file_remote(source, row_groups_read, row_groups_total);
super::convert::warn_spill_space(
source,
source.selected_input_bytes(selected_row_groups.as_ref())?,
options.spill_dir.as_deref(),
);
stage_input_pass0(source, selected_row_groups.as_ref(), row_groups_read);
let geom_idx = find_geometry_column(&input_schema).ok_or(ConvertError::NoGeometryColumn)?;
let geom_field = input_schema.field(geom_idx).clone();
let acc_cols = validate_cluster_schema(&input_schema, options)?;
validate_coalesce_schema(&input_schema, options)?;
Ok(Preflight {
options: resolved.clone(),
input_schema,
crs,
renames,
geom_idx,
geom_field,
acc_cols,
bbox_units,
bound_filter,
selected_row_groups,
row_groups_total,
row_groups_read,
})
}
pub(crate) fn convert_streaming_strategy(
source: &ConvertSource,
output_path: &Path,
options: &ConvertOptions,
strategy: Pass2Strategy,
) -> Result<ConvertReport, ConvertError> {
let start = Instant::now();
let mut peak_rss_mib = 0.0f64;
if options.sort_key.is_some() && options.class_ranking.is_some() {
return Err(ConvertError::RankingConflict);
}
let Preflight {
options: resolved_options,
input_schema,
crs,
renames,
geom_idx,
geom_field,
acc_cols,
bbox_units,
bound_filter,
selected_row_groups,
row_groups_total,
row_groups_read,
} = convert_preflight(source, options)?;
let options = &resolved_options;
let t_pass1 = Instant::now();
let Pass1Output {
mut features,
areas,
provenance: ranking_provenance,
acc_values,
coalesce: coalesce_scratch,
num_rows,
skipped_rows,
geom_bytes,
} = run_pass1(
source,
&input_schema,
geom_idx,
options,
&acc_cols,
selected_row_groups.as_ref(),
bbox_units.as_ref(),
bound_filter.as_ref(),
)?;
if skipped_rows > 0 {
log::warn!(
"skipping {skipped_rows} of {num_rows} input rows with a null, \
empty, or non-finite geometry"
);
}
let num_features = features.len();
let antimeridian_suspect_features = features
.iter()
.filter(|f| super::convert::bbox_antimeridian_suspect(&f.bbox, crs))
.count();
super::convert::warn_antimeridian_suspects(antimeridian_suspect_features);
log::info!("[convert] scan complete: {num_features} feature(s) from {num_rows} row(s)");
log::debug!(
"[profile] pass1 stream+scan: {:.2}s",
t_pass1.elapsed().as_secs_f64()
);
log_phase_rss("pass1 scan", &mut peak_rss_mib);
let WinnerTables {
level_specs,
cluster_tables,
kinds,
coalesce_scratch,
min_levels,
counts,
carriers,
finest,
} = resolve_winner_tables(
&mut features,
acc_values,
areas,
coalesce_scratch,
num_rows,
crs,
options,
&mut peak_rss_mib,
)?;
let (emitted, mut skipped) = partition_emitted_levels(&level_specs, &counts);
if emitted.is_empty() {
return Err(ConvertError::NoData);
}
warn_plan_skipped_levels(&skipped, num_features, emitted[0].gsd, emitted[0].zoom);
let LevelWriter {
mut writer,
source_schema,
cluster_schema,
out_schema,
non_geom_cols,
} = create_level_writer(
output_path,
&input_schema,
geom_idx,
&geom_field,
&emitted,
crs,
ranking_provenance,
&renames,
options,
)?;
log_phase_rss("pre-pass2 (winner tables freed)", &mut peak_rss_mib);
let t_pass2 = Instant::now();
let coalesce_tables =
build_pass2_coalesce_tables(coalesce_scratch.as_ref(), &emitted, finest, crs, options);
let duplicating = matches!(options.mode, Mode::Duplicating);
let cascade_chains = build_cascade_chains(&emitted, finest, duplicating, options);
let ctxs = build_level_ctxs(
&emitted,
options,
&LevelCtxInputs {
source_schema: &source_schema,
cluster_schema: &cluster_schema,
out_schema: &out_schema,
non_geom_cols: &non_geom_cols,
geom_idx,
min_levels: &min_levels,
acc_cols: &acc_cols,
kinds: kinds.as_deref(),
cluster_tables: cluster_tables.as_ref(),
coalesce_tables: &coalesce_tables,
cascade_chains: &cascade_chains,
carriers: &carriers,
crs,
finest,
duplicating,
},
);
let hints: Vec<usize> = emitted.iter().map(|e| e.hint).collect();
let validation_skips_before = validation_skip_count();
let in_flight_batches = resolve_and_log_in_flight_batches(options.in_flight_batches);
let level_stats = run_pass2_levels(
&mut writer,
&ctxs,
&hints,
source,
options,
selected_row_groups.as_ref(),
in_flight_batches,
&out_schema,
num_rows,
geom_bytes,
strategy,
)?;
log_validation_skips(validation_skips_before);
let mut level_reports = Vec::with_capacity(emitted.len());
for (e, (outcome, rows, vertices)) in emitted.iter().zip(level_stats) {
record_level_outcome(
outcome,
SkippedLevelReport {
planned_level: e.orig as usize,
gsd: e.gsd,
zoom: e.zoom,
},
e.hint,
rows,
vertices,
&mut level_reports,
&mut skipped,
);
}
skipped.sort_by_key(|s| s.planned_level);
if level_reports.is_empty() {
return Err(ConvertError::NoData);
}
log::debug!(
"[profile] pass2 total: {:.2}s",
t_pass2.elapsed().as_secs_f64()
);
log_phase_rss("pass2 (output sink)", &mut peak_rss_mib);
let t_finish = Instant::now();
let meta = writer.finish()?;
log::debug!(
"[profile] writer.finish: {:.2}s",
t_finish.elapsed().as_secs_f64()
);
log_phase_rss("writer.finish", &mut peak_rss_mib);
log::info!("[rss] convert peak: {peak_rss_mib:.0} MiB");
fill_level_bytes(output_path, &meta, &mut level_reports)?;
let total_rows: usize = level_reports.iter().map(|l| l.feature_count).sum();
let total_vertices: usize = level_reports.iter().map(|l| l.vertex_count).sum();
let total_compressed_bytes: i64 = level_reports.iter().map(|l| l.compressed_bytes).sum();
Ok(ConvertReport {
mode: options.mode,
levels: level_reports,
skipped_empty_levels: skipped,
input_features: num_features,
total_rows,
total_vertices,
total_compressed_bytes,
row_groups_total,
row_groups_read,
antimeridian_suspect_features,
duration_secs: start.elapsed().as_secs_f64(),
remote_fetch: super::convert::log_remote_fetch(source),
})
}
struct RoadCandidate {
idx: usize,
ranking: ClassRanking,
found: HashSet<&'static str>,
keys: Vec<Option<f64>>,
groups: Vec<u32>,
interner: GroupInterner,
}
enum RankPlan {
ExplicitSort {
idx: usize,
name: String,
},
ExplicitClass {
idx: usize,
ranking: ClassRanking,
},
Auto {
roads: Vec<RoadCandidate>,
confidence: Option<(usize, String)>,
},
SizeFallback,
}
fn build_rank_plan(schema: &Schema, options: &ConvertOptions) -> Result<RankPlan, ConvertError> {
if let Some(name) = &options.sort_key {
let idx = schema
.index_of(name)
.map_err(|_| ConvertError::SortKeyColumnMissing { name: name.clone() })?;
return Ok(RankPlan::ExplicitSort {
idx,
name: name.clone(),
});
}
if let Some(cr) = &options.class_ranking {
let idx =
schema
.index_of(&cr.column)
.map_err(|_| ConvertError::ClassRankColumnMissing {
name: cr.column.clone(),
})?;
let dt = schema.field(idx).data_type();
if !matches!(dt, DataType::Utf8 | DataType::LargeUtf8) {
return Err(ConvertError::ClassRankColumnNotString {
name: cr.column.clone(),
data_type: format!("{dt:?}"),
});
}
return Ok(RankPlan::ExplicitClass {
idx,
ranking: cr.clone(),
});
}
if !options.no_auto_rank {
let roads: Vec<RoadCandidate> = schema
.fields()
.iter()
.enumerate()
.filter(|(_, f)| {
let lname = f.name().to_ascii_lowercase();
(lname == "road_class" || lname == "class")
&& matches!(f.data_type(), DataType::Utf8 | DataType::LargeUtf8)
})
.map(|(idx, f)| RoadCandidate {
idx,
ranking: overture_road_ranking(f.name().clone()),
found: HashSet::new(),
keys: Vec::new(),
groups: Vec::new(),
interner: GroupInterner::default(),
})
.collect();
let confidence = schema
.fields()
.iter()
.enumerate()
.find(|(_, f)| {
f.name().eq_ignore_ascii_case("confidence")
&& matches!(f.data_type(), DataType::Float32 | DataType::Float64)
})
.map(|(idx, f)| (idx, f.name().clone()));
if !roads.is_empty() || confidence.is_some() {
return Ok(RankPlan::Auto { roads, confidence });
}
}
Ok(RankPlan::SizeFallback)
}
fn scan_road_vocab(col: &dyn Array, found: &mut HashSet<&'static str>) {
use arrow_array::cast::AsArray;
if found.len() >= ROAD_VOCAB_MIN_DISTINCT {
return;
}
let vocab: HashSet<&'static str> = KNOWN_ROAD_CLASSES.iter().copied().collect();
macro_rules! scan {
($arr:expr) => {{
let a = $arr;
for i in 0..a.len() {
if a.is_null(i) {
continue;
}
if let Some(&hit) = vocab.get(a.value(i)) {
found.insert(hit);
if found.len() >= ROAD_VOCAB_MIN_DISTINCT {
return;
}
}
}
}};
}
match col.data_type() {
DataType::Utf8 => scan!(col.as_string::<i32>()),
DataType::LargeUtf8 => scan!(col.as_string::<i64>()),
_ => {}
}
}
struct CoalesceScratch {
rows: Vec<usize>,
geoms: Vec<Geometry<f64>>,
sort_keys: Vec<Option<f64>>,
groups: Option<Vec<u32>>,
}
impl CoalesceScratch {
fn inputs(&self) -> Vec<CoalesceInput<'_>> {
(0..self.rows.len())
.map(|i| CoalesceInput {
index: self.rows[i],
geom: &self.geoms[i],
sort_key: self.sort_keys[i],
group: self.groups.as_ref().map_or(0, |g| g[i]),
})
.collect()
}
}
struct Pass1Output {
features: Vec<AssignFeature>,
areas: Vec<f32>,
provenance: RankingProvenance,
acc_values: Vec<Vec<Option<f64>>>,
coalesce: Option<CoalesceScratch>,
num_rows: usize,
skipped_rows: usize,
geom_bytes: u64,
}
fn pass1_projection(
geom_idx: usize,
plan: &RankPlan,
acc_cols: &[usize],
filter: Option<&super::filter::BoundFilter>,
ladder_col: Option<usize>,
) -> Vec<usize> {
let mut cols: Vec<usize> = vec![geom_idx];
cols.extend(ladder_col);
if let Some(f) = filter {
cols.extend(f.columns().iter().copied());
}
match plan {
RankPlan::ExplicitSort { idx, .. } | RankPlan::ExplicitClass { idx, .. } => cols.push(*idx),
RankPlan::Auto { roads, confidence } => {
cols.extend(roads.iter().map(|r| r.idx));
if let Some((idx, _)) = confidence {
cols.push(*idx);
}
}
RankPlan::SizeFallback => {}
}
cols.extend(acc_cols.iter().copied());
cols.sort_unstable();
cols.dedup();
cols
}
fn apply_entry_levels(
options: &ConvertOptions,
ladder_values: &[Option<f64>],
num_rows: usize,
features: &mut [AssignFeature],
) -> Result<(), ConvertError> {
if options.entry_zoom.is_none() {
return Ok(());
}
debug_assert_eq!(ladder_values.len(), num_rows);
let level_specs = options.levels.resolve(options.gsd_base)?;
if let Some(entry) = super::convert::resolve_entry_levels(options, ladder_values, &level_specs)?
{
for f in features.iter_mut() {
f.entry_level = entry.get(f.index).copied().flatten();
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn run_pass1(
source: &ConvertSource,
input_schema: &Schema,
geom_idx: usize,
options: &ConvertOptions,
acc_cols: &[usize],
row_groups: Option<&RowGroupSelection>,
bbox_units: Option<&[f64; 4]>,
filter: Option<&super::filter::BoundFilter>,
) -> Result<Pass1Output, ConvertError> {
let mut plan = build_rank_plan(input_schema, options)?;
let ladder_col = options
.entry_zoom
.as_ref()
.map(|spec| {
input_schema.index_of(&spec.column).map_err(|_| {
ConvertError::InvalidConfig(format!(
"entry-zoom column {:?} not found in the input schema",
spec.column
))
})
})
.transpose()?;
let cols = pass1_projection(geom_idx, &plan, acc_cols, filter, ladder_col);
let proj = |orig: usize| cols.binary_search(&orig).expect("projected column");
let reader = source.open_stream(&ReadPlan {
batch_size: options.read_batch_size.max(1),
projection: Some(&cols),
row_groups,
})?;
let mut features: Vec<AssignFeature> = Vec::new();
let want_areas = accumulator_enabled(options);
let mut areas: Vec<f32> = Vec::new();
let mut num_rows = 0usize;
let mut geom_bytes = 0u64;
let mut skipped_rows = 0usize;
let mut point_count = 0usize;
let mut explicit_keys: Vec<Option<f64>> = Vec::new();
let mut confidence_keys: Vec<Option<f64>> = Vec::new();
let mut acc_values: Vec<Vec<Option<f64>>> = vec![Vec::new(); acc_cols.len()];
let mut ladder_values: Vec<Option<f64>> = Vec::new();
let mut geoms_buf: Vec<Option<Geometry<f64>>> = Vec::new();
let collect_lines = options.coalesce_lines;
let mut line_rows: Vec<usize> = Vec::new();
let mut line_feat_pos: Vec<usize> = Vec::new();
let mut line_geoms: Vec<Geometry<f64>> = Vec::new();
let mut explicit_groups: Vec<u32> = Vec::new();
let mut explicit_interner = GroupInterner::default();
for batch in reader {
let batch = batch?;
let gcol_idx = proj(geom_idx);
let schema = batch.schema();
let gfield = schema.field(gcol_idx);
let garr = from_arrow_array(batch.column(gcol_idx).as_ref(), gfield)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
geoms_buf.clear();
extract_geometries_opt_from_array(garr.as_ref(), &mut geoms_buf)?;
let filter_mask: Option<Vec<Option<bool>>> = filter.map(|f| f.eval_mask(&batch, &proj));
let base = num_rows;
let mut kept_row = vec![false; geoms_buf.len()];
for (i, gopt) in geoms_buf.iter().enumerate() {
if let Some(mask) = &filter_mask {
if mask[i] != Some(true) {
continue;
}
}
let Some(g) = gopt.as_ref() else {
skipped_rows += 1;
continue;
};
let Some((kind, fbbox)) = scan_feature(g) else {
skipped_rows += 1;
continue;
};
if let Some(bb) = bbox_units {
if !super::convert::bboxes_intersect(&fbbox, bb) {
continue;
}
}
if matches!(kind, FeatureKind::Point) {
point_count += 1;
}
if collect_lines && matches!(kind, FeatureKind::Line) {
line_rows.push(base + i);
line_feat_pos.push(features.len());
line_geoms.push(g.clone());
}
kept_row[i] = true;
if want_areas {
areas.push(polygon_area_f32(g));
}
features.push(AssignFeature {
index: base + i,
bbox: fbbox,
kind,
sort_key: None, entry_level: None,
});
}
num_rows += geoms_buf.len();
geom_bytes += batch.column(gcol_idx).get_array_memory_size() as u64;
match &mut plan {
RankPlan::ExplicitSort { idx, .. } => {
explicit_keys.extend(extract_sort_keys(batch.column(proj(*idx)).as_ref()));
}
RankPlan::ExplicitClass { idx, ranking } => {
let col = batch.column(proj(*idx));
explicit_keys.extend(extract_class_ranks(col.as_ref(), ranking)?);
if collect_lines {
explicit_interner.extend(col.as_ref(), &mut explicit_groups);
}
}
RankPlan::Auto { roads, confidence } => {
for cand in roads.iter_mut() {
let col = batch.column(proj(cand.idx));
scan_road_vocab(col.as_ref(), &mut cand.found);
cand.keys
.extend(extract_class_ranks(col.as_ref(), &cand.ranking)?);
if collect_lines {
cand.interner.extend(col.as_ref(), &mut cand.groups);
}
}
if let Some((idx, _)) = confidence {
confidence_keys.extend(extract_sort_keys(batch.column(proj(*idx)).as_ref()));
}
}
RankPlan::SizeFallback => {}
}
for (s, &idx) in acc_cols.iter().enumerate() {
acc_values[s].extend(extract_sort_keys(batch.column(proj(idx)).as_ref()));
}
if let Some(idx) = ladder_col {
let keys = extract_sort_keys(batch.column(proj(idx)).as_ref());
ladder_values.extend(
keys.into_iter()
.zip(&kept_row)
.map(|(k, keep)| if *keep { k } else { None }),
);
}
}
let (keys, provenance, all_groups) = resolve_ranking_tier(
plan,
explicit_keys,
confidence_keys,
explicit_groups,
collect_lines,
features.len(),
point_count,
);
if let Some(keys) = keys {
debug_assert_eq!(keys.len(), num_rows);
for f in features.iter_mut() {
f.sort_key = keys[f.index];
}
}
apply_entry_levels(options, &ladder_values, num_rows, &mut features)?;
let coalesce = collect_lines.then(|| CoalesceScratch {
sort_keys: line_feat_pos
.iter()
.map(|&p| features[p].sort_key)
.collect(),
groups: all_groups.map(|g| line_rows.iter().map(|&r| g[r]).collect()),
rows: line_rows,
geoms: line_geoms,
});
Ok(Pass1Output {
features,
areas,
provenance,
acc_values,
coalesce,
num_rows,
skipped_rows,
geom_bytes,
})
}
#[derive(Default)]
pub(super) struct Pass2Timers {
read: AtomicU64,
decode: AtomicU64,
simplify: AtomicU64,
build: AtomicU64,
}
impl Pass2Timers {
fn add(cell: &AtomicU64, start: Instant) {
cell.fetch_add(start.elapsed().as_nanos() as u64, Ordering::Relaxed);
}
pub(super) fn add_dur(cell: &AtomicU64, dur: Duration) {
cell.fetch_add(dur.as_nanos() as u64, Ordering::Relaxed);
}
fn secs(cell: &AtomicU64) -> f64 {
Duration::from_nanos(cell.load(Ordering::Relaxed)).as_secs_f64()
}
pub(super) fn read_cell(&self) -> &AtomicU64 {
&self.read
}
pub(super) fn log_engine_summary(&self, total_secs: f64, rows: usize) {
let read_s = Self::secs(&self.read);
let decode_s = Self::secs(&self.decode);
let simplify_s = Self::secs(&self.simplify);
let build_s = Self::secs(&self.build);
log::debug!(
"[profile] pass2 engine ({rows} rows): wall={total_secs:.2}s \
read={read_s:.2}s decode={decode_s:.2}s simplify={simplify_s:.2}s \
build={build_s:.2}s (stage sums are core-seconds, overlap wall)"
);
}
}
pub(super) struct LevelStreamCtx<'a> {
source_schema: &'a Schema,
cluster_schema: &'a Schema,
out_schema: &'a Schema,
non_geom_cols: &'a [usize],
geom_idx: usize,
min_levels: &'a [u8],
orig_level: u8,
duplicating: bool,
verbatim: bool,
gsd_m: f64,
repr: Representation,
crs: Crs,
simplify: &'a SimplifyOptions,
cluster_enabled: bool,
cluster_table: Option<&'a std::collections::HashMap<usize, ClusterEntry>>,
acc_cols: &'a [usize],
coalesce_enabled: bool,
kinds: Option<&'a [FeatureKind]>,
coalesce_table: Option<&'a CoalesceTable>,
cascade_chain: &'a [CascadeStep],
carriers: &'a [usize],
}
impl LevelStreamCtx<'_> {
#[inline]
fn is_member(&self, g: usize) -> bool {
let ml = self.min_levels[g];
if self.duplicating {
ml <= self.orig_level || is_carrier(self.carriers, g)
} else {
ml == self.orig_level
}
}
#[inline]
fn is_carrier_row(&self, g: usize) -> bool {
self.duplicating && self.min_levels[g] > self.orig_level && is_carrier(self.carriers, g)
}
}
impl LevelStreamCtx<'_> {
pub(super) fn is_cascading_duplicating(&self) -> bool {
self.duplicating && self.simplify.cascade
}
}
#[allow(clippy::too_many_arguments)]
fn write_level_streaming(
writer: &mut OverviewWriter<File>,
level_idx: usize,
hint: usize,
source: &ConvertSource,
read_batch_size: usize,
in_flight: usize,
row_groups: Option<&RowGroupSelection>,
ctx: &LevelStreamCtx<'_>,
) -> Result<(LevelWriteOutcome, usize, usize), ConvertError> {
let rows = Cell::new(0usize);
let vertices = Cell::new(0usize);
let recv_wait_ns = Cell::new(0u64);
let timers = Pass2Timers::default();
let fallbacks_before = full_resolution_fallback_count();
let t_level = Instant::now();
struct Processed {
batch: RecordBatch,
verts: usize,
}
let timers = &timers;
let outcome = scoped_pipe(
in_flight,
|tx: &Sender<Processed>| -> Result<(), ConvertError> {
let mut reader = source.open_stream(&ReadPlan {
batch_size: read_batch_size.max(1),
projection: None,
row_groups,
})?;
let mut row_offset = 0usize;
let mut last_progress = Instant::now();
loop {
if last_progress.elapsed().as_secs() >= 10 {
last_progress = Instant::now();
log::info!(
"[convert] level {level_idx}: {row_offset} input \
row(s) scanned",
);
}
let t_read = Instant::now();
let batch = match reader.next() {
None => return Ok(()),
Some(Err(e)) => return Err(e.into()),
Some(Ok(b)) => b,
};
Pass2Timers::add(&timers.read, t_read);
let offset = row_offset;
row_offset += batch.num_rows();
match process_level_batch(&batch, offset, ctx, timers)? {
None => continue, Some((out, verts)) => {
if tx.send(Processed { batch: out, verts }).is_err() {
return Ok(());
}
}
}
}
},
|rx: Receiver<Processed>| -> Result<LevelWriteOutcome, ConvertError> {
let batches = std::iter::from_fn(|| {
let t_wait = Instant::now();
match rx.recv() {
Ok(msg) => {
recv_wait_ns.set(recv_wait_ns.get() + t_wait.elapsed().as_nanos() as u64);
rows.set(rows.get() + msg.batch.num_rows());
vertices.set(vertices.get() + msg.verts);
Some(msg.batch)
}
Err(_) => {
recv_wait_ns.set(recv_wait_ns.get() + t_wait.elapsed().as_nanos() as u64);
None
}
}
});
Ok(writer.write_level(level_idx, Some(hint), batches)?)
},
)?;
let total = t_level.elapsed().as_secs_f64();
let read_s = Pass2Timers::secs(&timers.read);
let decode_s = Pass2Timers::secs(&timers.decode);
let simplify_s = Pass2Timers::secs(&timers.simplify);
let build_s = Pass2Timers::secs(&timers.build);
let writer_busy = total - Duration::from_nanos(recv_wait_ns.get()).as_secs_f64();
log::debug!(
"[profile] level {} ({}, {} rows): total={:.2}s read={:.2}s decode={:.2}s \
simplify={:.2}s build={:.2}s writer_busy={:.2}s (read/decode/simplify/build \
overlap the writer)",
level_idx,
if ctx.verbatim { "verbatim" } else { "simplify" },
rows.get(),
total,
read_s,
decode_s,
simplify_s,
build_s,
writer_busy,
);
let fallbacks = full_resolution_fallback_count() - fallbacks_before;
if fallbacks > 0 {
log::debug!(
"[profile] level {level_idx}: {fallbacks} feature(s) kept at full \
resolution (invalid RDP candidate after all epsilon retries)"
);
}
Ok((outcome, rows.get(), vertices.get()))
}
pub(super) fn process_level_batch(
batch: &RecordBatch,
row_offset: usize,
ctx: &LevelStreamCtx<'_>,
timers: &Pass2Timers,
) -> Result<Option<(RecordBatch, usize)>, ConvertError> {
let n = batch.num_rows();
let t_decode = Instant::now();
let selected: Vec<usize> = (0..n)
.filter(|&i| {
let g = row_offset + i;
if let Some(table) = ctx.coalesce_table {
if ctx.kinds.expect("kinds present when coalescing")[g] == FeatureKind::Line {
return table.contains_key(&g);
}
}
ctx.is_member(g)
})
.collect();
if selected.is_empty() {
return Ok(None);
}
let take_idx = UInt32Array::from(selected.iter().map(|&i| i as u32).collect::<Vec<_>>());
let geom_taken = take(batch.column(ctx.geom_idx).as_ref(), &take_idx, None)?;
let schema = batch.schema();
let gfield = schema.field(ctx.geom_idx);
let garr = from_arrow_array(geom_taken.as_ref(), gfield)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(selected.len());
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
Pass2Timers::add(&timers.decode, t_decode);
let t_simplify = Instant::now();
let mut kept_idx: Vec<usize> = Vec::with_capacity(selected.len());
let mut verts = 0usize;
let kept_geoms: Vec<Geometry<f64>> = if ctx.verbatim {
for (g, &i) in geoms.iter().zip(&selected) {
verts += count_vertices(g);
kept_idx.push(i);
}
geoms
} else {
let simplified: Vec<Simplified> = geoms
.par_iter()
.zip(&selected)
.map(|(g, &i)| {
if let Some((merged, _)) = ctx.coalesce_table.and_then(|t| t.get(&(row_offset + i)))
{
Simplified::Keep(merged.clone())
} else if ctx.is_carrier_row(row_offset + i) {
carrier_square(g, ctx.gsd_m, ctx.crs, ctx.simplify)
.map_or(Simplified::Dropped, Simplified::Keep)
} else if !ctx.cascade_chain.is_empty() {
simplify_cascade(g, ctx.cascade_chain, ctx.crs, ctx.simplify)
} else {
simplify_step(g, ctx.gsd_m, ctx.crs, ctx.simplify, ctx.repr)
}
})
.collect();
let mut out = Vec::with_capacity(selected.len());
for (s, &i) in simplified.into_iter().zip(&selected) {
match s {
Simplified::Keep(s) => {
verts += count_vertices(&s);
kept_idx.push(i);
out.push(s);
}
Simplified::Dropped => {}
}
}
if out.is_empty() {
Pass2Timers::add(&timers.simplify, t_simplify);
return Ok(None);
}
out
};
Pass2Timers::add(&timers.simplify, t_simplify);
let t_build = Instant::now();
let out_batch = assemble_level_batch(batch, row_offset, ctx, &kept_idx, &kept_geoms)?;
Pass2Timers::add(&timers.build, t_build);
Ok(Some((out_batch, verts)))
}
fn assemble_level_batch(
batch: &RecordBatch,
row_offset: usize,
ctx: &LevelStreamCtx<'_>,
kept_idx: &[usize],
kept_geoms: &[Geometry<f64>],
) -> Result<RecordBatch, ConvertError> {
let mut out_batch = build_level_batch(
ctx.source_schema,
batch,
ctx.non_geom_cols,
ctx.geom_idx,
kept_idx,
kept_geoms,
)?;
if ctx.cluster_enabled || ctx.coalesce_enabled {
let globals: Vec<usize> = kept_idx.iter().map(|&i| row_offset + i).collect();
if ctx.cluster_enabled {
out_batch = apply_cluster_columns(
out_batch,
ctx.cluster_schema,
&globals,
ctx.cluster_table,
ctx.acc_cols,
)?;
}
if ctx.coalesce_enabled {
out_batch =
apply_coalesced_count(out_batch, ctx.out_schema, &globals, ctx.coalesce_table)?;
}
}
Ok(out_batch)
}
pub(super) fn process_batch_cascade(
batch: &RecordBatch,
row_offset: usize,
ctxs: &[LevelStreamCtx<'_>],
timers: &Pass2Timers,
) -> Result<Vec<Option<(RecordBatch, usize)>>, ConvertError> {
let Some(finest) = ctxs.last() else {
return Ok(Vec::new());
};
debug_assert!(ctxs.iter().all(|c| c.duplicating && !c.verbatim));
debug_assert!(ctxs
.iter()
.enumerate()
.all(|(li, c)| c.cascade_chain.len() == ctxs.len() - li
&& c.cascade_chain.last()
== Some(&CascadeStep {
gsd_meters: c.gsd_m,
repr: c.repr,
})));
debug_assert!(ctxs
.iter()
.all(|c| c.coalesce_table.is_some() == finest.coalesce_table.is_some()));
let n = batch.num_rows();
let t_decode = Instant::now();
let mut pos_of_row: Vec<u32> = vec![u32::MAX; n];
let mut selected: Vec<usize> = Vec::with_capacity(n);
for (i, pos) in pos_of_row.iter_mut().enumerate() {
let g = row_offset + i;
if finest.coalesce_table.is_some()
&& finest.kinds.expect("kinds present when coalescing")[g] == FeatureKind::Line
{
continue;
}
if finest.min_levels[g] <= finest.orig_level
|| ctxs.iter().any(|c| is_carrier(c.carriers, g))
{
*pos = u32::try_from(selected.len()).expect("batch rows fit in u32");
selected.push(i);
}
}
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(selected.len());
if !selected.is_empty() {
let take_idx = UInt32Array::from(selected.iter().map(|&i| i as u32).collect::<Vec<_>>());
let geom_taken = take(batch.column(finest.geom_idx).as_ref(), &take_idx, None)?;
let schema = batch.schema();
let gfield = schema.field(finest.geom_idx);
let garr = from_arrow_array(geom_taken.as_ref(), gfield)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
}
Pass2Timers::add(&timers.decode, t_decode);
let has_band_upto: Vec<bool> = {
let mut v = Vec::with_capacity(ctxs.len());
let mut any = false;
for c in ctxs.iter() {
any = any || c.repr != Representation::Geometry;
v.push(any);
}
v
};
let t_simplify = Instant::now();
let folds: Vec<Vec<Simplified>> = geoms
.par_iter()
.zip(&selected)
.map(|(g, &i)| {
let ml = finest.min_levels[row_offset + i];
let mut out: Vec<Simplified> = Vec::with_capacity(ctxs.len());
let mut current: Option<Geometry<f64>> = None;
let mut alive = true;
for (li, ctx) in ctxs.iter().enumerate().rev() {
if ml > ctx.orig_level {
break; }
if !alive && !has_band_upto[li] {
break; }
let step = if !alive && ctx.repr == Representation::Geometry {
Simplified::Dropped
} else {
let input = if alive {
current.as_ref().unwrap_or(g)
} else {
g };
simplify_step(input, ctx.gsd_m, ctx.crs, ctx.simplify, ctx.repr)
};
match step {
Simplified::Keep(s) => {
out.push(Simplified::Keep(s.clone()));
current = Some(s);
alive = true;
}
Simplified::Dropped => {
out.push(Simplified::Dropped);
alive = false;
}
}
}
out
})
.collect();
Pass2Timers::add(&timers.simplify, t_simplify);
let t_build = Instant::now();
let results: Vec<Result<Option<(RecordBatch, usize)>, ConvertError>> = ctxs
.par_iter()
.enumerate()
.map(|(li, ctx)| {
let depth = ctxs.len() - 1 - li;
let mut kept_idx: Vec<usize> = Vec::new();
let mut kept_geoms: Vec<Geometry<f64>> = Vec::new();
let mut verts = 0usize;
for (i, &pos) in pos_of_row.iter().enumerate() {
let g = row_offset + i;
if let Some(table) = ctx.coalesce_table {
if ctx.kinds.expect("kinds present when coalescing")[g] == FeatureKind::Line {
if let Some((merged, _)) = table.get(&g) {
verts += count_vertices(merged);
kept_idx.push(i);
kept_geoms.push(merged.clone());
}
continue;
}
}
if ctx.min_levels[g] <= ctx.orig_level {
debug_assert_ne!(pos, u32::MAX, "member row missing from cascade superset");
if let Some(Simplified::Keep(s)) = folds[pos as usize].get(depth) {
verts += count_vertices(s);
kept_idx.push(i);
kept_geoms.push(s.clone());
}
} else if ctx.is_carrier_row(g) {
debug_assert_ne!(pos, u32::MAX, "carrier row missing from cascade superset");
if let Some(sq) =
carrier_square(&geoms[pos as usize], ctx.gsd_m, ctx.crs, ctx.simplify)
{
verts += count_vertices(&sq);
kept_idx.push(i);
kept_geoms.push(sq);
}
}
}
if kept_idx.is_empty() {
return Ok(None);
}
let out_batch = assemble_level_batch(batch, row_offset, ctx, &kept_idx, &kept_geoms)?;
Ok(Some((out_batch, verts)))
})
.collect();
Pass2Timers::add(&timers.build, t_build);
let mut per_level = Vec::with_capacity(results.len());
for res in results {
per_level.push(res?);
}
Ok(per_level)
}