mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
perf(lens): faster trace opens and list pages at scale (#45228)
* perf(lens): read trace spans in larger batches Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * perf(lens): evaluate the trace list page once per identity lookup Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> * fix(lens): queue trace reads instead of rejecting past 8 in flight Co-Authored-By: Ishaan Jaffer <155045088+ishaan-berri@users.noreply.github.com> --------- Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
ef694bc584
commit
b970e412d9
6 changed files with 83 additions and 51 deletions
|
|
@ -42,5 +42,6 @@ prettyplease = "0.2"
|
|||
|
||||
[dev-dependencies]
|
||||
rstest.workspace = true
|
||||
tokio = { workspace = true, features = ["test-util"] }
|
||||
wiremock.workspace = true
|
||||
uuid.workspace = true
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ use litellm_traces_clickhouse::InsertTable;
|
|||
use serde_json::Value;
|
||||
use std::{
|
||||
collections::BTreeMap,
|
||||
future::Future,
|
||||
sync::{
|
||||
Arc,
|
||||
atomic::{AtomicBool, Ordering},
|
||||
|
|
@ -33,6 +34,7 @@ use std::{
|
|||
time::Duration,
|
||||
};
|
||||
pub use storage::Storage;
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
#[allow(
|
||||
dead_code,
|
||||
|
|
@ -46,7 +48,8 @@ pub use storage::Storage;
|
|||
pub mod wire {
|
||||
include!(concat!(env!("OUT_DIR"), "/wire.rs"));
|
||||
}
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
const READ_QUEUE_WAIT: Duration = Duration::from_secs(10);
|
||||
|
||||
pub struct State {
|
||||
pub credentials: Arc<auth::Credentials>,
|
||||
|
|
@ -80,6 +83,15 @@ impl State {
|
|||
}
|
||||
}
|
||||
|
||||
async fn wait_for_read_slot<P>(
|
||||
acquire: impl Future<Output = Result<P, tokio::sync::AcquireError>>,
|
||||
) -> Result<P, Error> {
|
||||
tokio::time::timeout(READ_QUEUE_WAIT, acquire)
|
||||
.await
|
||||
.map_err(|_| Error::Unavailable)?
|
||||
.map_err(|_| Error::Unavailable)
|
||||
}
|
||||
|
||||
pub fn router(state: Arc<State>) -> Router {
|
||||
let public = Router::new()
|
||||
.route("/health/live", get(|| async { StatusCode::OK }))
|
||||
|
|
@ -125,10 +137,7 @@ async fn receipt(
|
|||
) -> Result<Json<Value>, Error> {
|
||||
let tenant = state.credentials.tenant(&headers)?;
|
||||
state.require_storage()?;
|
||||
let _permit = state
|
||||
.read_slots
|
||||
.try_acquire()
|
||||
.map_err(|_| Error::Unavailable)?;
|
||||
let _permit = wait_for_read_slot(state.read_slots.acquire()).await?;
|
||||
let body = tokio::time::timeout(Duration::from_secs(5), to_bytes(body, 64 * 1024))
|
||||
.await
|
||||
.map_err(|_| Error::Unavailable)?
|
||||
|
|
@ -206,11 +215,7 @@ async fn read(
|
|||
) -> Result<Json<Value>, Error> {
|
||||
auth::authorize_service(&headers, &state.service_token)?;
|
||||
state.require_storage()?;
|
||||
let permit = state
|
||||
.read_slots
|
||||
.clone()
|
||||
.try_acquire_owned()
|
||||
.map_err(|_| Error::Unavailable)?;
|
||||
let permit = wait_for_read_slot(state.read_slots.clone().acquire_owned()).await?;
|
||||
let body = tokio::time::timeout(Duration::from_secs(10), to_bytes(body, 1024 * 1024))
|
||||
.await
|
||||
.map_err(|_| Error::Unavailable)?
|
||||
|
|
@ -279,3 +284,39 @@ pub async fn provision(state: Arc<State>) {
|
|||
tokio::time::sleep(Duration::from_secs(10)).await;
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{Error, READ_QUEUE_WAIT, wait_for_read_slot};
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
#[tokio::test]
|
||||
async fn ninth_read_waits_for_a_permit_and_succeeds() {
|
||||
let slots = Arc::new(Semaphore::new(8));
|
||||
let permits = (0..8)
|
||||
.map(|_| slots.clone().try_acquire_owned().expect("available permit"))
|
||||
.collect::<Vec<_>>();
|
||||
let waiting_slots = slots.clone();
|
||||
let waiting =
|
||||
tokio::spawn(async move { wait_for_read_slot(waiting_slots.acquire_owned()).await });
|
||||
|
||||
tokio::task::yield_now().await;
|
||||
assert!(!waiting.is_finished());
|
||||
drop(permits);
|
||||
assert!(waiting.await.expect("joined read").is_ok());
|
||||
}
|
||||
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn read_queue_timeout_returns_unavailable() {
|
||||
let slots = Arc::new(Semaphore::new(0));
|
||||
let waiting = tokio::spawn(wait_for_read_slot(slots.acquire_owned()));
|
||||
|
||||
tokio::task::yield_now().await;
|
||||
tokio::time::advance(READ_QUEUE_WAIT).await;
|
||||
assert!(matches!(
|
||||
waiting.await.expect("joined read"),
|
||||
Err(Error::Unavailable)
|
||||
));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ SELECT TraceId AS trace_id,
|
|||
ifNull(any(RootName), '') AS name, any(ServiceName) AS service,
|
||||
ifNull(any(RootInput), '') AS input_preview, ifNull(any(RootStatus), '') AS status,
|
||||
toUnixTimestamp64Milli(min(StartTs)) AS start_ms,
|
||||
min(StartTs) AS trace_start, max(EndTs) AS trace_end,
|
||||
dateDiff('millisecond', min(StartTs), max(EndTs)) AS duration_ms,
|
||||
sum(SpanCount) AS span_count,
|
||||
sum(AgentCount) AS agent_invocations,
|
||||
|
|
@ -27,7 +26,7 @@ HAVING min(StartTs) >= fromUnixTimestamp64Milli({start_ms:Int64})
|
|||
ORDER BY start_ms DESC, trace_ref DESC
|
||||
LIMIT {limit:UInt32}
|
||||
)
|
||||
SELECT page.* EXCEPT (trace_start, trace_end),
|
||||
SELECT page.*,
|
||||
identities.agent_names AS agent_names, identities.agent_count AS agent_count,
|
||||
identities.frameworks AS frameworks
|
||||
FROM page
|
||||
|
|
@ -37,10 +36,8 @@ LEFT JOIN (
|
|||
arraySort(groupUniqArrayIf(toString(Framework), Framework != '')) AS frameworks,
|
||||
uniqExactIf(if(AgentName = '', SpanName, AgentName), ObservationType = 'agent') AS agent_count
|
||||
FROM otel_traces
|
||||
WHERE Timestamp >= (SELECT min(trace_start) FROM page)
|
||||
AND Timestamp <= (SELECT max(trace_end) FROM page)
|
||||
WHERE Timestamp >= fromUnixTimestamp64Milli({start_ms:Int64})
|
||||
AND TraceId IN (SELECT trace_id FROM page)
|
||||
AND (TeamId, ApiKeyHash, TraceId) IN (SELECT team_id, api_key_hash, trace_id FROM page)
|
||||
GROUP BY TeamId, ApiKeyHash, TraceId
|
||||
) AS identities
|
||||
ON page.team_id = identities.TeamId AND page.api_key_hash = identities.ApiKeyHash
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@
|
|||
use std::{future::Future, marker::PhantomData};
|
||||
|
||||
use litellm_http::Client;
|
||||
use litellm_storage_clickhouse::{Query, fetch};
|
||||
use litellm_storage_clickhouse::{Query, ReadLimits, fetch};
|
||||
use litellm_traces::query::named as contracts;
|
||||
use litellm_traces_cache::{MAX_GRAPH_BYTES, MAX_GRAPH_SPANS, StoreError};
|
||||
use serde::{Serialize, de::DeserializeOwned};
|
||||
|
|
@ -14,7 +14,12 @@ use crate::{
|
|||
query::named::{SpendByResponseIdsParams, SpendByResponseIdsRow, TraceSpansRow},
|
||||
};
|
||||
|
||||
const PAGE_SIZE: u32 = 256;
|
||||
const PAGE_SIZE: u32 = 8192;
|
||||
const SPAN_BATCH_READ_LIMITS: ReadLimits = ReadLimits {
|
||||
result_rows: PAGE_SIZE as u64,
|
||||
response_bytes: 16 * 1024 * 1024,
|
||||
..litellm_storage_clickhouse::READ_LIMITS
|
||||
};
|
||||
|
||||
#[derive(Default)]
|
||||
struct ReadBudget {
|
||||
|
|
@ -67,6 +72,7 @@ struct Paged<K>(PhantomData<K>);
|
|||
impl<K: Keyset> Query for Paged<K> {
|
||||
type Params = Batch<K>;
|
||||
type Row = K::Row;
|
||||
const READ_LIMITS: ReadLimits = SPAN_BATCH_READ_LIMITS;
|
||||
const SQL: &'static str = K::SQL;
|
||||
}
|
||||
|
||||
|
|
@ -316,8 +322,8 @@ mod tests {
|
|||
}
|
||||
|
||||
#[rstest]
|
||||
#[case::fits(1000, PAGE_SIZE, &[256, 256, 256, 256])]
|
||||
#[case::uniform_large_rows(1000, 100, &[256, 128, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64, 64])]
|
||||
#[case::fits(1000, PAGE_SIZE, &[8192])]
|
||||
#[case::uniform_large_rows(200, 100, &[8192, 4096, 2048, 1024, 512, 256, 128, 64, 64, 64, 64])]
|
||||
#[tokio::test]
|
||||
async fn a_rejected_page_size_is_not_retried(
|
||||
#[case] total: u32,
|
||||
|
|
@ -334,6 +340,13 @@ mod tests {
|
|||
assert_eq!(table.requests.lock().unwrap().as_slice(), requests);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn span_batches_use_larger_read_limits() {
|
||||
let limits = <Paged<Numbers> as Query>::READ_LIMITS;
|
||||
assert_eq!(limits.result_rows, 8192);
|
||||
assert_eq!(limits.response_bytes, 16 * 1024 * 1024);
|
||||
}
|
||||
|
||||
#[rstest]
|
||||
#[tokio::test]
|
||||
async fn a_single_oversized_row_fails_the_read() {
|
||||
|
|
@ -346,7 +359,9 @@ mod tests {
|
|||
assert!(matches!(result, Err(StoreError::TooLarge)), "{result:?}");
|
||||
assert_eq!(
|
||||
table.requests.lock().unwrap().as_slice(),
|
||||
&[256, 128, 64, 32, 16, 8, 4, 2, 1]
|
||||
&[
|
||||
8192, 4096, 2048, 1024, 512, 256, 128, 64, 32, 16, 8, 4, 2, 1
|
||||
]
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ pub const DATABASE: &str = "trace_test";
|
|||
|
||||
pub struct SeededDatabase {
|
||||
pub database: ClickHouseDatabase,
|
||||
#[allow(dead_code)] // dead_code: also shared with reads.rs, which uses a direct storage reader
|
||||
pub readers: QueryReaders,
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,9 +3,7 @@ use std::collections::BTreeMap;
|
|||
use litellm_http::Client;
|
||||
use litellm_traces::query::named::ReadAccessParams;
|
||||
use litellm_traces_cache::{ReadError, TraceReader};
|
||||
use litellm_traces_clickhouse::{
|
||||
ClickHouseTraces, Connection, InsertTable, QueryScope, insert_rows,
|
||||
};
|
||||
use litellm_traces_clickhouse::{ClickHouseTraces, Connection, InsertTable, insert_rows};
|
||||
use rstest::rstest;
|
||||
use serde_json::json;
|
||||
|
||||
|
|
@ -83,10 +81,7 @@ async fn list_costs_match_each_run_when_response_ids_are_reused(
|
|||
.collect(),
|
||||
)
|
||||
.await?;
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection);
|
||||
let access = ReadAccessParams {
|
||||
all_teams: false,
|
||||
|
|
@ -212,10 +207,7 @@ async fn large_runs_remain_complete_under_default_reader_limits(
|
|||
.collect::<Vec<_>>();
|
||||
insert_rows(client, &writer, DATABASE, InsertTable::SpendLogs, costs).await?;
|
||||
}
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection);
|
||||
let access = ReadAccessParams {
|
||||
all_teams: false,
|
||||
|
|
@ -365,10 +357,7 @@ async fn cursor_pages_keep_a_tenant_scoped_snapshot_when_more_spans_arrive(
|
|||
) -> TestResult {
|
||||
let fixture = seeded_database?;
|
||||
let client = &fixture.database.client;
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection.clone());
|
||||
let access = ReadAccessParams {
|
||||
all_teams: true,
|
||||
|
|
@ -528,10 +517,7 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals(
|
|||
) -> TestResult {
|
||||
let fixture = seeded_database?;
|
||||
let client = &fixture.database.client;
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection.clone());
|
||||
let access = ReadAccessParams {
|
||||
all_teams: true,
|
||||
|
|
@ -557,10 +543,7 @@ async fn an_oversized_span_keeps_the_run_list_available_with_partial_totals(
|
|||
("TraceId".into(), json!(run.trace_id)),
|
||||
("SpanId".into(), json!("oversized-child")),
|
||||
("ParentSpanId".into(), json!("0101010101010101")),
|
||||
(
|
||||
"SpanName".into(),
|
||||
json!("x".repeat(litellm_storage_clickhouse::READ_LIMITS.response_bytes + 1)),
|
||||
),
|
||||
("SpanName".into(), json!("x".repeat(16 * 1024 * 1024 + 1))),
|
||||
("ObservationType".into(), json!("tool")),
|
||||
("TeamId".into(), json!("team-a")),
|
||||
("ApiKeyHash".into(), json!("key-a")),
|
||||
|
|
@ -699,10 +682,7 @@ async fn assigned_call_ids_require_shared_ownership_through_detail_and_batch_rea
|
|||
.collect(),
|
||||
)
|
||||
.await?;
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection);
|
||||
let access = ReadAccessParams {
|
||||
all_teams: false,
|
||||
|
|
@ -815,10 +795,7 @@ async fn native_cost_correlation_survives_session_grouping_and_excludes_other_ow
|
|||
.collect(),
|
||||
)
|
||||
.await?;
|
||||
let connection = fixture
|
||||
.readers
|
||||
.connection(client, &QueryScope::All, "fixture-secret")
|
||||
.await?;
|
||||
let connection = Connection::reader(&fixture.database.url, DATABASE)?;
|
||||
let (reader, store) = make_reader(client, connection);
|
||||
let access = ReadAccessParams {
|
||||
all_teams: true,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue