Skip to content
Draft
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
15 changes: 7 additions & 8 deletions quickwit/quickwit-common/src/tracing_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,14 +80,13 @@ pub fn extract_context(metadata: &MetadataMap) -> Context {
global::get_text_map_propagator(|propagator| propagator.extract(&extractor))
}

/// Extracts a W3C trace context from incoming gRPC request metadata and
/// installs it as the parent of the currently-active tracing span. Use this
/// at the entry of a gRPC handler that is itself wrapped in a
/// `#[tracing::instrument]` so the handler's span is stitched into the
/// caller's trace.
pub fn set_current_span_parent_from_metadata(metadata: &MetadataMap) {
let parent_context = extract_context(metadata);
let _ = Span::current().set_parent(parent_context);
/// Records an attribute on the currently-active span's OpenTelemetry span.
///
/// Unlike a `tracing` field (`#[instrument(fields(...))]` or `Span::record`), this does
/// not propagate to log events, so it can carry verbose values (e.g. a query AST)
/// without bloating logs. No-op when no OpenTelemetry layer is installed.
pub fn record_current_span_attribute(key: &'static str, value: impl Into<opentelemetry::Value>) {
Span::current().set_attribute(key, value);
}

/// Tonic interceptor that injects the active span's W3C trace context into
Expand Down
3 changes: 2 additions & 1 deletion quickwit/quickwit-search/src/fetch_docs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ use tantivy::schema::document::CompactDocValue;
use tantivy::schema::{Document as DocumentTrait, Field, TantivyDocument, Value};
use tantivy::snippet::SnippetGenerator;
use tantivy::{ReloadPolicy, Score, Searcher, Term};
use tracing::{Instrument, error};
use tracing::{Instrument, error, instrument};

use crate::leaf::open_index_with_caches;
use crate::service::SearcherContext;
Expand Down Expand Up @@ -153,6 +153,7 @@ struct Document {
}

/// Fetching docs from a specific split.
#[instrument(skip_all, fields(split_id = split.split_id, num_docs = global_doc_addrs.len()))]
async fn fetch_docs_in_split(
searcher_context: Arc<SearcherContext>,
mut global_doc_addrs: Vec<GlobalDocAddress>,
Expand Down
116 changes: 98 additions & 18 deletions quickwit/quickwit-search/src/leaf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ use crate::metrics::{
};
use crate::root::is_metadata_count_request_with_ast;
use crate::search_permit_provider::{
SearchPermit, SearchPermitFuture, compute_initial_memory_allocation,
BlockReasonHandle, SearchPermit, SearchPermitFuture, compute_initial_memory_allocation,
};
use crate::service::{SearcherContext, deserialize_doc_mapper};
use crate::{QuickwitAggregations, SearchError};
Expand Down Expand Up @@ -300,7 +300,6 @@ async fn run_cancellable(
///
/// Returns whether the query is provably empty in this split (i.e. `on_absent` fired and
/// warmup was short-circuited).
#[instrument(skip_all)]
pub(crate) async fn warmup(
searcher: &Searcher,
warmup_info: &WarmupInfo,
Expand Down Expand Up @@ -779,7 +778,27 @@ async fn leaf_search_single_split(
HitSet::empty(),
);
};
warmup(&searcher, &warmup_info, &record_absence).await?
// `warmup_mb` is the total warmed into the byte-range cache; `downloaded_mb` is
// the subset actually fetched from object storage (cache hits excluded). Both are
// recorded after warmup, before the search adds to either.
const BYTES_PER_MB: f64 = 1_000_000.0;
let warmup_span = info_span!(
"warmup",
warmup_mb = tracing::field::Empty,
downloaded_mb = tracing::field::Empty
);
let provably_empty = warmup(&searcher, &warmup_info, &record_absence)
.instrument(warmup_span.clone())
.await?;
warmup_span.record(
"warmup_mb",
byte_range_cache.get_num_bytes() as f64 / BYTES_PER_MB,
);
warmup_span.record(
"downloaded_mb",
download_counters.snapshot().0 as f64 / BYTES_PER_MB,
);
provably_empty
};
let warmup_end = Instant::now();
let warmup_duration: Duration = warmup_end.duration_since(warmup_start);
Expand Down Expand Up @@ -841,7 +860,10 @@ async fn leaf_search_single_split(
}

let split_num_docs = split.num_docs;
let span = info_span!("tantivy_search");
// Spans the CPU-pool queue wait: created here (right after warmup) and closed when the
// closure starts executing on a pool thread. `tantivy_search` is created inside the
// closure so it covers only the CPU execution, not this wait.
let cpu_wait_span = info_span!("waiting_for_cpu_pool");

let split_clone = split.clone();

Expand All @@ -852,9 +874,12 @@ async fn leaf_search_single_split(
let search_request_and_result: Option<(SearchRequest, LeafSearchResponse)> =
crate::search_thread_pool()
.run_cpu_intensive(move || {
// The CPU-pool queue wait ends as this closure starts executing.
drop(cpu_wait_span);
leaf_search_state_guard.set_state(SplitSearchState::Cpu);
let cpu_start = Instant::now();
let cpu_thread_pool_wait_microsecs = cpu_start.duration_since(warmup_end);
let span = info_span!("tantivy_search");
let _span_guard = span.enter();
// Our search execution has been scheduled, let's check if we can improve the
// request based on the results of the preceding searches
Expand Down Expand Up @@ -1545,7 +1570,7 @@ pub async fn multi_index_leaf_search(
let searcher_context = searcher_context.clone();
let search_request = search_request.clone();

leaf_request_futures.spawn({
leaf_request_futures.spawn(
async move {
let storage = storage_resolver.resolve(&index_uri).await?;
single_doc_mapping_leaf_search(
Expand All @@ -1555,10 +1580,10 @@ pub async fn multi_index_leaf_search(
leaf_search_request_ref.split_offsets,
doc_mapper,
)
.in_current_span()
.await
}
});
.instrument(Span::current()),
);
}

// Creates a collector which merges responses into one
Expand Down Expand Up @@ -1723,12 +1748,15 @@ async fn run_offloaded_search_tasks(
}],
};
let invoker = lambda_invoker.clone();
lambda_tasks_joinset.spawn(async move {
(
batch_split_ids,
invoker.invoke_leaf_search(leaf_request).await,
)
});
lambda_tasks_joinset.spawn(
async move {
(
batch_split_ids,
invoker.invoke_leaf_search(leaf_request).await,
)
}
.in_current_span(),
);
}

while let Some(join_res) = lambda_tasks_joinset.join_next().await {
Expand Down Expand Up @@ -2006,6 +2034,42 @@ pub async fn single_doc_mapping_leaf_search(
Ok(leaf_search_response)
}

/// Records the permit block reason onto the `waiting_for_leaf_search_split_semaphore`
/// span when the wait ends.
///
/// Implemented as a drop guard (rather than reading after `.await`) so the reason is
/// recorded even when the wait future is cancelled before the permit is granted — the
/// long waits that hit a search deadline are exactly the ones that never get granted.
struct WaitBlockReasonRecorder {
wait_span: Span,
block_reason: BlockReasonHandle,
started_at: Instant,
}

impl WaitBlockReasonRecorder {
fn new(wait_span: Span, block_reason: BlockReasonHandle) -> Self {
Self {
wait_span,
block_reason,
started_at: Instant::now(),
}
}
}

impl Drop for WaitBlockReasonRecorder {
fn drop(&mut self) {
// Skip near-instant grants: a sub-millisecond wait isn't worth attributing and
// keeps the span uncluttered.
const MIN_WAIT_FOR_BLOCK_ATTRIBUTION: Duration = Duration::from_millis(1);
if self.started_at.elapsed() < MIN_WAIT_FOR_BLOCK_ATTRIBUTION {
return;
}
if let Some(block_reason) = self.block_reason.get() {
self.wait_span.record("blocked_on", block_reason.as_str());
}
}
}

async fn run_local_search_tasks(
local_search_tasks: Vec<LocalSearchTask>,
index_storage: Arc<dyn Storage + 'static>,
Expand All @@ -2021,9 +2085,26 @@ async fn run_local_search_tasks(
search_permit_future,
} in local_search_tasks
{
let leaf_split_search_permit = search_permit_future
.instrument(info_span!("waiting_for_leaf_search_split_semaphore"))
.await;
// Per-split span covering both the permit wait and the search, so each split is a
// single subtree (wait + warmup + tantivy) rather than flat siblings.
let split_span = info_span!(
"leaf_search_single_split_wrapper",
split_id = split.split_id,
num_docs = split.num_docs
);
let wait_span = info_span!(
parent: &split_span,
"waiting_for_leaf_search_split_semaphore",
blocked_on = tracing::field::Empty
);
// Records the block reason on `wait_span` when the wait ends — whether the permit
// is granted or the wait is cancelled on a search deadline (the important, long
// waits are exactly the ones that get cancelled before being granted).
let _block_reason_recorder = WaitBlockReasonRecorder::new(
wait_span.clone(),
search_permit_future.block_reason_handle(),
);
let leaf_split_search_permit = search_permit_future.instrument(wait_span).await;

// We run simplify search request again: as we push split into the merge collector,
// we may have discovered that we won't find any better candidates for top hits in this
Expand All @@ -2045,7 +2126,7 @@ async fn run_local_search_tasks(
split.clone(),
leaf_split_search_permit,
)
.in_current_span(),
.instrument(split_span),
);
task_id_to_split_id_map.insert(handle.id(), split_id);
}
Expand Down Expand Up @@ -2174,7 +2255,6 @@ struct LeafSearchContext {
split_filter: Arc<RwLock<CanSplitDoBetter>>,
}

#[instrument(skip_all, fields(split_id = split.split_id, num_docs = split.num_docs))]
async fn leaf_search_single_split_wrapper(
request: SearchRequest,
ctx: Arc<LeafSearchContext>,
Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-search/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,7 @@ pub async fn list_all_splits(
}

/// Extract the list of relevant splits for a given request.
#[tracing::instrument(skip_all, fields(num_indexes = index_uids.len()))]
pub async fn list_relevant_splits(
index_uids: Vec<IndexUid>,
start_timestamp: Option<i64>,
Expand Down
26 changes: 23 additions & 3 deletions quickwit/quickwit-search/src/root.rs
Original file line number Diff line number Diff line change
Expand Up @@ -789,7 +789,7 @@ fn compute_root_resource_stats(

/// If this method fails for some splits, a partial search response is returned, with the list of
/// faulty splits in the failed_splits field.
#[instrument(level = "debug", skip_all)]
#[instrument(skip_all)]
pub(crate) async fn search_partial_hits_phase(
searcher_context: &SearcherContext,
indexes_metas_for_leaf_search: &IndexesMetasForLeafSearch,
Expand Down Expand Up @@ -1249,6 +1249,7 @@ async fn refine_and_list_matches(
}

/// Fetches the list of splits and their metadata from the metastore
#[instrument(skip_all)]
async fn plan_splits_for_root_search(
search_request: &mut SearchRequest,
metastore: &mut MetastoreServiceClient,
Expand Down Expand Up @@ -1311,8 +1312,9 @@ pub async fn root_search(
let num_docs: usize = split_metadatas.iter().map(|split| split.num_docs).sum();
let num_splits = split_metadatas.len();

// It would have been nice to add those in the context of the trace span,
// but with our current logging setting, it makes logs too verbose.
// These would bloat the logs if added as tracing span fields (they propagate to
// every log event in the span), so they go on a one-shot `info!` for logs and, via
// `record_current_span_attribute`, on the OpenTelemetry span only for the trace.
info!(
query_ast = search_request.query_ast.as_str(),
agg = search_request.aggregation_request(),
Expand All @@ -1322,6 +1324,24 @@ pub async fn root_search(
num_splits = num_splits,
"root_search"
);
quickwit_common::tracing_utils::record_current_span_attribute(
"query_ast",
search_request.query_ast.clone(),
);
quickwit_common::tracing_utils::record_current_span_attribute("num_docs", num_docs as i64);
quickwit_common::tracing_utils::record_current_span_attribute("num_splits", num_splits as i64);
if let Some(start_timestamp) = search_request.start_timestamp {
quickwit_common::tracing_utils::record_current_span_attribute(
"start_timestamp",
start_timestamp,
);
}
if let Some(end_timestamp) = search_request.end_timestamp {
quickwit_common::tracing_utils::record_current_span_attribute(
"end_timestamp",
end_timestamp,
);
}

if let Some(max_total_split_searches) = searcher_context.searcher_config.max_splits_per_search
&& max_total_split_searches < num_splits
Expand Down
Loading
Loading