Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/bin/stream/bootstrap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -428,6 +428,7 @@ pub(crate) async fn run_bootstrap(
by_database: vec![boot_db.series()],
..stage_gauges(&StageCounters {
emitter: Some(&stats),
db_name: &|_| boot_db.database.clone(),
oracle: [None, Some(&oracle_stats)],
bootstrap: Some(&progress),
bootstrap_attempt: attempt,
Expand Down
70 changes: 70 additions & 0 deletions src/bin/stream/metrics_publish.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use walshadow::boundary_hold::BoundaryHoldStats;
use walshadow::ch_emitter::EmitterStats;
use walshadow::config::ConfigResolver;
use walshadow::metrics::{DbSeries, MetricsRegistry, MetricsSnapshot};
use walshadow::pipeline::row_ledger::TableRowCounts;
use walshadow::pos::{
Drain, EmitterAck, FilterDispatched, Floor, Pos, ShadowReplay, SourceReceived,
};
Expand Down Expand Up @@ -416,6 +417,8 @@ pub(crate) async fn populate_pipeline_metrics(
/// counters sum into one series
pub(crate) struct StageCounters<'a> {
pub(crate) emitter: Option<&'a walshadow::ch_emitter::EmitterStats>,
/// `database=` label for a source database oid
pub(crate) db_name: &'a (dyn Fn(u32) -> String + Sync),
/// Live first, bootstrap's second. Only one of the pair is ever serving
pub(crate) oracle: [Option<&'a walshadow::oracle::OracleStats>; 2],
pub(crate) bootstrap: Option<&'a BootstrapProgress>,
Expand All @@ -432,12 +435,26 @@ pub(crate) fn emitter_counts<const N: usize>(
picks.map(|pick| stats.map_or(0, |s| pick(s).load(Ordering::Relaxed)))
}

/// Stop conditions beyond the crossing and shadow views
pub(crate) fn stage_gauges(v: &StageCounters<'_>) -> MetricsSnapshot {
stage_gauges_on(v, MetricsSnapshot::default())
}

const UNMAPPED_STATUS_ROWS: usize = 20;

pub(crate) fn stage_gauges_on(v: &StageCounters<'_>, base: MetricsSnapshot) -> MetricsSnapshot {
let (proc_cpu, proc_rss, proc_threads) = read_process_stats();
let by_table = |pick: fn(&EmitterStats) -> &TableRowCounts| {
v.emitter.map_or_else(Vec::new, |s| {
pick(s)
.snapshot()
.into_iter()
.map(|(key, rows)| ((v.db_name)(key.db_oid), key.rel.to_string(), rows))
.collect::<Vec<_>>()
})
};
let mut unmapped_rows_by_table = by_table(|s| &s.unmapped_rows_by_table);
unmapped_rows_by_table.truncate(UNMAPPED_STATUS_ROWS);
let emitter = |pick: fn(&EmitterStats) -> &AtomicU64| -> u64 {
v.emitter.map_or(0, |s| pick(s).load(Ordering::Relaxed))
};
Expand Down Expand Up @@ -538,6 +555,17 @@ pub(crate) fn stage_gauges_on(v: &StageCounters<'_>, base: MetricsSnapshot) -> M
emitter_xacts_total: emitter(|s| &s.xacts_committed),
emitter_unsupported_relations: emitter(|s| &s.unsupported_relations),
emitter_deletes_discarded: emitter(|s| &s.deletes_discarded),
unmapped_rows_by_table,
rows_inserted_by_table: by_table(|s| &s.rows_inserted_by_table),
cells_inserted_by_type: v.emitter.map_or_else(Vec::new, |s| {
s.cells_inserted_by_type
.snapshot()
.into_iter()
.map(|((ty, enc), n)| (ty, enc, n))
.collect()
}),
emitter_retries_attempted_total: emitter(|s| &s.retries_attempted),
emitter_reconnects_total: emitter(|s| &s.reconnects),
oracle_local_columns_total: emitter(|s| &s.oracle_local_columns),
oracle_blocks_total: oracle(|s| &s.blocks),
oracle_rows_total: oracle(|s| &s.rows),
Expand Down Expand Up @@ -643,6 +671,7 @@ mod tests {
boundary_hold: &boundary,
by_database: Vec::new(),
counters: StageCounters {
db_name: &|_| String::new(),
emitter: Some(&emitter),
oracle: [None, None],
bootstrap: None,
Expand All @@ -667,11 +696,50 @@ mod tests {
}
}

/// The publisher names each counted table by resolving its `db_oid`, so a
/// ledger with rows has to reach the snapshot under the database's name
#[test]
fn per_table_counters_resolve_their_database_name() {
use walshadow::schema::{RelName, TableKey};
let emitter = EmitterStats::default();
let orders = TableKey::new(5, RelName::new("public", "orders"));
let noise = TableKey::new(9, RelName::new("public", "noise"));
emitter
.rows_inserted_by_table
.counter(&orders)
.fetch_add(7, Ordering::Relaxed);
emitter
.unmapped_rows_by_table
.counter(&noise)
.fetch_add(3, Ordering::Relaxed);
let snap = stage_gauges(&StageCounters {
db_name: &|oid| match oid {
5 => "app".to_owned(),
_ => format!("db{oid}"),
},
emitter: Some(&emitter),
oracle: [None, None],
bootstrap: None,
bootstrap_attempt: 0,
uptime_secs: 0,
});
assert_eq!(
snap.rows_inserted_by_table,
vec![("app".to_owned(), "public.orders".to_owned(), 7)],
);
assert_eq!(
snap.unmapped_rows_by_table,
vec![("db9".to_owned(), "public.noise".to_owned(), 3)],
"an oid the config does not name still has to be reportable",
);
}

#[test]
fn archive_stage_metrics_refresh_without_resetting_recovery_state() {
use walshadow::record::WAL_SEG_SIZE;
let emitter = EmitterStats::default();
let counters = StageCounters {
db_name: &|_| String::new(),
emitter: Some(&emitter),
oracle: [None, None],
bootstrap: None,
Expand Down Expand Up @@ -726,6 +794,7 @@ mod tests {
boot_oracle.rows.fetch_add(7, Ordering::Relaxed);

let snap = stage_gauges(&StageCounters {
db_name: &|_| String::new(),
emitter: None,
oracle: [Some(&live_oracle), Some(&boot_oracle)],
bootstrap: None,
Expand All @@ -737,6 +806,7 @@ mod tests {
assert_eq!(snap.bootstrap_attempt, 2);

let boot_only = stage_gauges(&StageCounters {
db_name: &|_| String::new(),
emitter: None,
oracle: [None, Some(&boot_oracle)],
bootstrap: None,
Expand Down
7 changes: 7 additions & 0 deletions src/bin/stream/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1895,6 +1895,13 @@ pub(crate) async fn run_session(
metrics_dbs.iter().map(DbMetricSources::series).collect(),
StageCounters {
emitter: emitter_stats,
db_name: &|oid| {
db_conns
.iter()
.find(|c| c.oid == oid)
.map(|c| c.name.clone())
.unwrap_or_default()
},
oracle: [oracle_stats, bootstrap_metrics.as_ref().map(|b| &*b.oracle)],
bootstrap: bootstrap_metrics.as_ref().map(|b| &b.progress),
bootstrap_attempt: 0,
Expand Down
1 change: 1 addition & 0 deletions src/bin/stream/tracing_setup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ pub(crate) fn init_tracing(
.with(filter)
.with(fmt_layer)
.with(otel_layer)
.with(walshadow::log_events::LogEventLayer)
.try_init();
provider
}
10 changes: 10 additions & 0 deletions src/emit/ch_emitter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1394,6 +1394,13 @@ pub(crate) enum ColumnEncoding {
}

impl ColumnEncoding {
pub(crate) fn label(&self) -> &'static str {
match self {
Self::Local => "local",
Self::Oracle { .. } => "oracle",
}
}

fn choose(att: Option<&RelAttr>, ast: &TypeAst) -> Self {
let Some(att) = att else {
// Predeclared post-ALTER column has only default cells
Expand Down Expand Up @@ -2551,6 +2558,9 @@ crate::atomic_stats! {
/// (memoised)
pub route_snapshots_mapped,
pub route_snapshots_unmapped,
pub unmapped_rows_by_table: crate::emit::pipeline::row_ledger::TableRowCounts,
pub rows_inserted_by_table: crate::emit::pipeline::row_ledger::TableRowCounts,
pub cells_inserted_by_type: crate::emit::pipeline::row_ledger::CellCounts,
/// Commit-resolve raw decode: records by verdict kind, per op
/// (`raw_decode_records_total{kind,op}`)
pub raw_decode_toast_ops: crate::decode::heap_decoder::OpCounters,
Expand Down
39 changes: 27 additions & 12 deletions src/emit/pipeline/batcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@

use std::collections::hash_map::Entry;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::time::Duration;

use clickhouse_c::Allocator;
Expand All @@ -32,7 +32,8 @@ use tokio::time::Instant;
use crate::config::ResolvedConfig;
use crate::decode::heap_decoder::{CommittedTuple, HeapOp};
use crate::emit::ch_emitter::{
Append, ColumnBuf, EmitterStats, OP_DELETE, OP_INSERT, OP_UPDATE, TableEncoder, TablePlan,
Append, ColumnBuf, ColumnEncoding, EmitterStats, OP_DELETE, OP_INSERT, OP_UPDATE, TableEncoder,
TablePlan,
};
use crate::emit::pipeline::{DEFAULT_PIPELINE_FLUSH, Fatal};
use crate::emit::route::RouteSnapshot;
Expand Down Expand Up @@ -67,6 +68,7 @@ pub struct RowChunk {
pub struct ColMeta {
pub name: String,
pub type_repr: String,
pub cells_inserted: Arc<AtomicU64>,
}

/// Immutable per-table block shape, shared by every batch of that table
Expand All @@ -78,16 +80,26 @@ pub struct BatchMeta {
/// synthetic ones (lsn, xid, commit_ts, delete marker when configured).
pub columns: Vec<ColMeta>,
pub schema_epoch: u64,
pub rows_inserted: Arc<AtomicU64>,
}

impl BatchMeta {
fn from_plan(plan: &TablePlan, table_key: TableKey, schema_epoch: u64) -> Self {
fn from_plan(
plan: &TablePlan,
table_key: TableKey,
schema_epoch: u64,
stats: &EmitterStats,
) -> Self {
let col_meta = |name: &String, type_repr: &String, encoding: &ColumnEncoding| ColMeta {
name: name.clone(),
type_repr: type_repr.clone(),
cells_inserted: stats
.cells_inserted_by_type
.counter(&(type_repr.clone(), encoding.label())),
};
let mut columns = Vec::with_capacity(plan.columns.len() + 4);
for c in &plan.columns {
columns.push(ColMeta {
name: c.name.clone(),
type_repr: c.type_repr.clone(),
});
columns.push(col_meta(&c.name, &c.type_repr, &c.encoding));
}
for synth in [
Some(&plan.synth_lsn),
Expand All @@ -98,12 +110,10 @@ impl BatchMeta {
.into_iter()
.flatten()
{
columns.push(ColMeta {
name: synth.name.clone(),
type_repr: synth.type_repr.clone(),
});
columns.push(col_meta(&synth.name, &synth.type_repr, &synth.encoding));
}
Self {
rows_inserted: stats.rows_inserted_by_table.counter(&table_key),
table_key,
insert_sql: plan.insert_sql.clone(),
columns,
Expand Down Expand Up @@ -333,7 +343,12 @@ async fn handle_row(
row.route.system_columns(),
)
.map_err(|e| e.to_string())?;
let meta = Arc::new(BatchMeta::from_plan(&plan, e.key().clone(), ctx.epoch));
let meta = Arc::new(BatchMeta::from_plan(
&plan,
e.key().clone(),
ctx.epoch,
ctx.stats,
));
let enc = TableEncoder::new(plan).map_err(|e| e.to_string())?;
e.insert(Table {
enc,
Expand Down
8 changes: 8 additions & 0 deletions src/emit/pipeline/inserter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,14 @@ impl Inserter {
self.stats
.rows_emitted
.fetch_add(batch.n_rows as u64, Ordering::Relaxed);
batch
.meta
.rows_inserted
.fetch_add(batch.n_rows as u64, Ordering::Relaxed);
for col in &batch.meta.columns {
col.cells_inserted
.fetch_add(batch.n_rows as u64, Ordering::Relaxed);
}
self.stats.blocks_sent.fetch_add(1, Ordering::Relaxed);
self.stats
.inserter_batches_in
Expand Down
1 change: 1 addition & 0 deletions src/emit/pipeline/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ pub mod plan_spool;
pub mod planner;
pub mod reorder;
pub mod resolver;
pub mod row_ledger;
pub mod tail;

use std::sync::Arc;
Expand Down
55 changes: 32 additions & 23 deletions src/emit/pipeline/reorder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ use crate::decode::visibility::{PgXactPatch, PgXactView, read_pg_xact};
use crate::emit::ch_ddl::DdlApplicator;
use crate::emit::ch_emitter::EmitterStats;
use crate::record::{Record, RecordSink, SinkError};
use crate::schema::{RelDescriptor, RelName, SchemaEvent};
use crate::schema::{RelDescriptor, RelName, SchemaEvent, TableKey};
use ahash::{HashMap, HashMapExt, HashSet, HashSetExt};
use tracing::Instrument;

Expand Down Expand Up @@ -1161,31 +1161,40 @@ impl<'a> ReorderRouteView<'a> {
impl PlanRouteView for ReorderRouteView<'_> {
fn route_for(&mut self, heap: &DescribedHeap) -> Option<Arc<RouteSnapshot>> {
let rel_name = &heap.descriptor.rel_name;
if let Some(r) = self.memo.get(rel_name) {
return r.clone();
}
let mapped = match self.overlay.get(rel_name) {
Some(o) => o.clone(),
None => self.mapping.as_ref().and_then(|m| m.get(rel_name)).cloned(),
let route = if let Some(r) = self.memo.get(rel_name) {
r.clone()
} else {
let mapped = match self.overlay.get(rel_name) {
Some(o) => o.clone(),
None => self.mapping.as_ref().and_then(|m| m.get(rel_name)).cloned(),
};
let route = mapped.map(|m| {
let rules = self
.config
.as_ref()
.map_or_else(Arc::default, |rc| rc.column_rules.clone());
let policy = self.row_policy.for_rel(self.config.as_deref(), rel_name);
RouteSnapshot::freeze(Arc::new(m), rules, policy)
});
let result = if route.is_none() {
self.stats
.unsupported_relations
.fetch_add(1, Ordering::Relaxed);
&self.stats.route_snapshots_unmapped
} else {
&self.stats.route_snapshots_mapped
};
result.fetch_add(1, Ordering::Relaxed);
self.memo.insert(rel_name.clone(), route.clone());
route
};
let route = mapped.map(|m| {
let rules = self
.config
.as_ref()
.map_or_else(Arc::default, |rc| rc.column_rules.clone());
let policy = self.row_policy.for_rel(self.config.as_deref(), rel_name);
RouteSnapshot::freeze(Arc::new(m), rules, policy)
});
let result = if route.is_none() {
if route.is_none() {
let key = TableKey::new(heap.descriptor.rfn.db_node, rel_name.clone());
self.stats
.unsupported_relations
.unmapped_rows_by_table
.counter(&key)
.fetch_add(1, Ordering::Relaxed);
&self.stats.route_snapshots_unmapped
} else {
&self.stats.route_snapshots_mapped
};
result.fetch_add(1, Ordering::Relaxed);
self.memo.insert(rel_name.clone(), route.clone());
}
route
}

Expand Down
Loading
Loading