Skip to content
1 change: 0 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 0 additions & 1 deletion benchmarks/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,6 @@ tempfile = "3"

[dev-dependencies]
criterion = "0.5"
sysinfo = "0.30"

[build-dependencies]
built = { version = "0.8", features = ["git2", "chrono"] }
Expand Down
259 changes: 137 additions & 122 deletions benchmarks/benches/broadcast_cache_scenarios.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,23 +5,20 @@ use datafusion::arrow::record_batch::RecordBatch;
use datafusion::common::Statistics;
use datafusion::common::tree_node::TreeNodeRecursion;
use datafusion::error::Result;
use datafusion::execution::memory_pool::{MemoryPool, PeakRecordingPool, UnboundedMemoryPool};
use datafusion::execution::runtime_env::RuntimeEnvBuilder;
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr};
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{
DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties,
};
use datafusion::prelude::SessionContext;
use datafusion::prelude::{SessionConfig, SessionContext};
use datafusion_distributed::BroadcastExec;
use futures::{StreamExt, stream};
use std::sync::{
Arc,
atomic::{AtomicBool, AtomicU64, Ordering},
};
use std::thread;
use std::sync::Arc;
use std::time::{Duration, Instant};
use sysinfo::{System, get_current_pid};
use tokio::runtime::Builder as RuntimeBuilder;

#[derive(Clone, Copy)]
Expand All @@ -47,56 +44,22 @@ struct Scenario {
consumers: Vec<ConsumerSpec>,
}

struct PeakRssSampler {
stop: Arc<AtomicBool>,
peak_kb: Arc<AtomicU64>,
handle: Option<thread::JoinHandle<()>>,
}

impl PeakRssSampler {
fn start(interval: Duration) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let peak_kb = Arc::new(AtomicU64::new(0));
let stop_clone = Arc::clone(&stop);
let peak_clone = Arc::clone(&peak_kb);
let handle = thread::spawn(move || {
let mut sys = System::new();
let pid = get_current_pid().expect("pid");
while !stop_clone.load(Ordering::Relaxed) {
sys.refresh_process(pid);
if let Some(proc) = sys.process(pid) {
let mem_kb = proc.memory();
peak_clone.fetch_max(mem_kb, Ordering::Relaxed);
}
thread::sleep(interval);
}
});
Self {
stop,
peak_kb,
handle: Some(handle),
}
}

fn stop(mut self) -> u64 {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
self.peak_kb.load(Ordering::Relaxed)
}
}

#[derive(Debug)]
struct SyntheticExec {
schema: SchemaRef,
partitions: usize,
batches: Arc<Vec<Arc<RecordBatch>>>,
interval: Option<Duration>,
properties: Arc<PlanProperties>,
}

impl SyntheticExec {
fn new(schema: SchemaRef, partitions: usize, batches: Arc<Vec<Arc<RecordBatch>>>) -> Self {
fn new(
schema: SchemaRef,
partitions: usize,
batches: Arc<Vec<Arc<RecordBatch>>>,
interval: Option<Duration>,
) -> Self {
let properties = Arc::new(PlanProperties::new(
EquivalenceProperties::new(Arc::clone(&schema)),
Partitioning::UnknownPartitioning(partitions),
Expand All @@ -107,6 +70,7 @@ impl SyntheticExec {
schema,
partitions,
batches,
interval,
properties,
}
}
Expand Down Expand Up @@ -158,12 +122,29 @@ impl ExecutionPlan for SyntheticExec {
let batches = Arc::clone(&self.batches);
let len = batches.len();

let stream = stream::iter((0..len).map(move |idx| {
let batch = &batches[idx];
Ok(batch.as_ref().clone())
}));

Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
match self.interval {
None => {
let stream = stream::iter((0..len).map(move |idx| {
let batch = &batches[idx];
Ok(batch.as_ref().clone())
}));
Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
}
Some(interval) => {
let stream = futures::stream::unfold(0usize, move |idx| {
let batches = Arc::clone(&batches);
async move {
if idx >= len {
return None;
}
tokio::time::sleep(interval).await;
let batch = batches[idx].as_ref().clone();
Some((Ok(batch), idx + 1))
}
});
Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
}
}
}

fn partition_statistics(&self, _partition: Option<usize>) -> Result<Arc<Statistics>> {
Expand Down Expand Up @@ -197,23 +178,15 @@ async fn consume_partition(

async fn run_scenario(
scenario: &Scenario,
schema: Arc<Schema>,
batches: Arc<Vec<Arc<RecordBatch>>>,
input: Arc<dyn ExecutionPlan>,
task_ctx: Arc<TaskContext>,
sample_rss: bool,
) -> Result<(Duration, u64)> {
let input: Arc<dyn ExecutionPlan> = Arc::new(SyntheticExec::new(
Arc::clone(&schema),
scenario.input_partitions,
batches,
));

recording_pool: Option<Arc<PeakRecordingPool>>,
) -> Result<(Duration, Option<(usize, usize)>)> {
let broadcast = Arc::new(BroadcastExec::new(
Arc::clone(&input),
scenario.consumer_tasks,
));

let sampler = sample_rss.then(|| PeakRssSampler::start(Duration::from_millis(25)));
let start = Instant::now();

let mut join_set = tokio::task::JoinSet::new();
Expand All @@ -232,8 +205,50 @@ async fn run_scenario(
}

let elapsed = start.elapsed();
let peak_kb = sampler.map(|s| s.stop()).unwrap_or(0);
Ok((elapsed, peak_kb))
let memory = recording_pool.map(|pool| (pool.peak_reserved(), pool.reserved()));
Ok((elapsed, memory))
}

fn memory_context() -> Result<(Arc<TaskContext>, Arc<PeakRecordingPool>)> {
let inner_pool: Arc<dyn MemoryPool> = Arc::new(UnboundedMemoryPool::default());
let recording_pool = Arc::new(PeakRecordingPool::new(inner_pool));
let runtime = RuntimeEnvBuilder::new()
.with_memory_pool(Arc::clone(&recording_pool) as Arc<dyn MemoryPool>)
.build()?;
let task_ctx =
SessionContext::new_with_config_rt(SessionConfig::new(), Arc::new(runtime)).task_ctx();
Ok((task_ctx, recording_pool))
}

fn synthetic_input(
scenario: &Scenario,
schema: &SchemaRef,
interval: Option<Duration>,
) -> Arc<dyn ExecutionPlan> {
let batches = (0..scenario.num_batches)
.map(|_| {
let array = UInt8Array::from(vec![0u8; scenario.rows_per_batch]);
let batch =
RecordBatch::try_new(Arc::clone(schema), vec![Arc::new(array)]).expect("batch");
Arc::new(batch)
})
.collect::<Vec<_>>();
Arc::new(SyntheticExec::new(
Arc::clone(schema),
scenario.input_partitions,
Arc::new(batches),
interval,
))
}

fn prebuilt_input(scenario: &Scenario, schema: &SchemaRef) -> Arc<dyn ExecutionPlan> {
synthetic_input(scenario, schema, None)
}

fn paced_input(scenario: &Scenario, schema: &SchemaRef) -> Arc<dyn ExecutionPlan> {
// Consumers are started before the first batch and receive a scheduling window between
// batches. This is intentionally outside the timing measurement.
synthetic_input(scenario, schema, Some(Duration::from_millis(1)))
}

fn all_fast_consumers(output_partitions: usize) -> Vec<ConsumerSpec> {
Expand Down Expand Up @@ -326,18 +341,8 @@ fn scenario_matrix() -> Vec<Scenario> {
scenarios
}

fn verbose_enabled() -> bool {
match std::env::var("BROADCAST_BENCH_VERBOSE") {
Ok(val) => {
let val = val.to_ascii_lowercase();
val == "1" || val == "true" || val == "yes"
}
Err(_) => false,
}
}

fn rss_enabled() -> bool {
match std::env::var("BROADCAST_BENCH_RSS") {
fn memory_enabled() -> bool {
match std::env::var("BROADCAST_BENCH_MEMORY") {
Ok(val) => {
let val = val.to_ascii_lowercase();
val == "1" || val == "true" || val == "yes"
Expand All @@ -353,6 +358,8 @@ fn runtime_threads() -> Option<usize> {
.filter(|threads| *threads > 0)
}

const DIAGNOSTIC_RUNS: usize = 5;

fn bench_broadcast_cache(c: &mut Criterion) {
let mut rt_builder = RuntimeBuilder::new_multi_thread();
if let Some(threads) = runtime_threads() {
Expand All @@ -362,8 +369,7 @@ fn bench_broadcast_cache(c: &mut Criterion) {

let mut group = c.benchmark_group("broadcast_cache_scenarios");
group.sample_size(10);
let verbose = verbose_enabled();
let sample_rss = rss_enabled();
let sample_memory = memory_enabled();
let task_ctx = SessionContext::new().task_ctx();

for scenario in scenario_matrix() {
Expand All @@ -372,60 +378,69 @@ fn bench_broadcast_cache(c: &mut Criterion) {
DataType::UInt8,
false,
)]));
let batches = (0..scenario.num_batches)
.map(|_| {
let data = vec![0u8; scenario.rows_per_batch];
let array = UInt8Array::from(data);
let batch = RecordBatch::try_new(Arc::clone(&schema), vec![Arc::new(array)])
.expect("batch");
Arc::new(batch)
})
.collect::<Vec<_>>();
let batches = Arc::new(batches);

let input = prebuilt_input(&scenario, &schema);
let mut memory_peaks = Vec::new();
let mut memory_residuals = Vec::new();
group.bench_function(BenchmarkId::new("scenario", scenario.name), |b| {
b.iter_custom(|iters| {
let mut total = Duration::ZERO;
let mut peaks = Vec::with_capacity(iters as usize);
for i in 0..iters {
let (elapsed, peak_kb) = rt
for _ in 0..iters {
let (elapsed, _) = rt
.block_on(run_scenario(
&scenario,
Arc::clone(&schema),
Arc::clone(&batches),
Arc::clone(&input),
Arc::clone(&task_ctx),
sample_rss,
None,
))
.expect("scenario");
if verbose || sample_rss {
eprintln!(
"scenario={} iter={} peak_rss_kb={} elapsed_ms={}",
scenario.name,
i,
peak_kb,
elapsed.as_millis()
);
}
peaks.push(peak_kb);
total += elapsed;
}
if sample_rss && !peaks.is_empty() {
peaks.sort_unstable();
let min = peaks[0];
let max = peaks[peaks.len() - 1];
let median = peaks[peaks.len() / 2];
eprintln!(
"scenario={} peak_rss_kb[min/median/max]={}/{}/{} runs={}",
scenario.name,
min,
median,
max,
peaks.len()
);
}
total
});
});

if sample_memory {
for _ in 0..DIAGNOSTIC_RUNS {
let (diagnostic_task_ctx, pool) = memory_context().expect("memory context");
let (_, memory) = rt
.block_on(run_scenario(
&scenario,
paced_input(&scenario, &schema),
diagnostic_task_ctx,
Some(pool),
))
.expect("diagnostic scenario");
if let Some((peak_reserved, residual_reserved)) = memory {
memory_peaks.push(peak_reserved);
memory_residuals.push(residual_reserved);
}
}
}

if sample_memory && !memory_peaks.is_empty() {
memory_peaks.sort_unstable();
memory_residuals.sort_unstable();
let peak_min = memory_peaks[0];
let peak_median = memory_peaks[memory_peaks.len() / 2];
let peak_max = memory_peaks[memory_peaks.len() - 1];
let residual_min = memory_residuals[0];
let residual_median = memory_residuals[memory_residuals.len() / 2];
let residual_max = memory_residuals[memory_residuals.len() - 1];
eprintln!(
concat!(
"scenario={} reserved_bytes[peak_min/median/max]={}/{}/{} ",
"end_min/median/max={}/{}/{} runs={}"
),
scenario.name,
peak_min,
peak_median,
peak_max,
residual_min,
residual_median,
residual_max,
memory_peaks.len()
);
}
}

group.finish();
Expand Down
Loading
Loading