fix(tracing): use one key-scoped trace rollup

This commit is contained in:
Yujong Lee 2026-09-30 13:18:56 -07:00
parent b354c738aa
commit b86acfbd7d
10 changed files with 41 additions and 70 deletions

View file

@ -24,7 +24,7 @@
<profile>litellm_traces_reader</profile>
<grants>
<query>GRANT SELECT ON litellm.otel_traces</query>
<query>GRANT SELECT ON litellm.agent_traces</query>
<query>GRANT SELECT ON litellm.agent_traces_by_key</query>
<query>GRANT SELECT ON litellm.spend_logs</query>
</grants>
</litellm_traces_reader>

View file

@ -1,6 +1,7 @@
CREATE TABLE IF NOT EXISTS {database}.agent_traces
CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key
(
TeamId LowCardinality(String),
ApiKeyHash String,
TraceId String,
StartTs SimpleAggregateFunction(min, DateTime64(9)),
EndTs SimpleAggregateFunction(max, DateTime64(9)),
@ -20,4 +21,4 @@ CREATE TABLE IF NOT EXISTS {database}.agent_traces
RequestIds SimpleAggregateFunction(groupArrayArray, Array(String))
)
ENGINE = AggregatingMergeTree
ORDER BY (TeamId, TraceId)
ORDER BY (TeamId, ApiKeyHash, TraceId)

View file

@ -1,6 +1,7 @@
CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_mv TO {database}.agent_traces AS
CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_by_key_mv
TO {database}.agent_traces_by_key AS
SELECT
TeamId, TraceId,
TeamId, ApiKeyHash, TraceId,
min(Timestamp) AS StartTs,
max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs,
any(ServiceName) AS ServiceName,
@ -18,4 +19,4 @@ SELECT
groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames,
groupArrayIf(LiteLLMRequestId, LiteLLMRequestId != '') AS RequestIds
FROM {database}.otel_traces
GROUP BY TeamId, TraceId
GROUP BY TeamId, ApiKeyHash, TraceId

View file

@ -1 +1 @@
ALTER TABLE {database}.agent_traces MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY
ALTER TABLE {database}.agent_traces_by_key MODIFY TTL toDateTime(StartTs) + INTERVAL {trace_retention_days} DAY

View file

@ -1,24 +0,0 @@
CREATE TABLE IF NOT EXISTS {database}.agent_traces_by_key
(
TeamId LowCardinality(String),
ApiKeyHash String,
TraceId String,
StartTs SimpleAggregateFunction(min, DateTime64(9)),
EndTs SimpleAggregateFunction(max, DateTime64(9)),
ServiceName SimpleAggregateFunction(any, LowCardinality(String)),
RootName SimpleAggregateFunction(anyLast, Nullable(String)),
RootInput SimpleAggregateFunction(anyLast, Nullable(String)),
RootStatus SimpleAggregateFunction(anyLast, Nullable(String)),
SpanCount SimpleAggregateFunction(sum, UInt64),
AgentCount SimpleAggregateFunction(sum, UInt64),
LlmCount SimpleAggregateFunction(sum, UInt64),
ToolCount SimpleAggregateFunction(sum, UInt64),
ErrorCount SimpleAggregateFunction(sum, UInt64),
InputTokens SimpleAggregateFunction(sum, UInt64),
OutputTokens SimpleAggregateFunction(sum, UInt64),
Models SimpleAggregateFunction(groupUniqArrayArray, Array(String)),
AgentNames SimpleAggregateFunction(groupUniqArrayArray, Array(String)),
RequestIds SimpleAggregateFunction(groupArrayArray, Array(String))
)
ENGINE = AggregatingMergeTree
ORDER BY (TeamId, ApiKeyHash, TraceId)

View file

@ -1,22 +0,0 @@
CREATE MATERIALIZED VIEW IF NOT EXISTS {database}.agent_traces_by_key_mv
TO {database}.agent_traces_by_key AS
SELECT
TeamId, ApiKeyHash, TraceId,
min(Timestamp) AS StartTs,
max(Timestamp + toIntervalNanosecond(Duration)) AS EndTs,
any(ServiceName) AS ServiceName,
anyLastIf(toNullable(SpanName), ParentSpanId = '') AS RootName,
anyLastIf(toNullable(InputPreview), ParentSpanId = '') AS RootInput,
anyLastIf(toNullable(StatusCode), ParentSpanId = '') AS RootStatus,
count() AS SpanCount,
countIf(ObservationType = 'agent') AS AgentCount,
countIf(ObservationType = 'llm') AS LlmCount,
countIf(ObservationType = 'tool') AS ToolCount,
countIf(StatusCode = 'STATUS_CODE_ERROR') AS ErrorCount,
sum(InputTokens) AS InputTokens,
sum(OutputTokens) AS OutputTokens,
groupUniqArrayIf(toString(Model), Model != '') AS Models,
groupUniqArrayIf(SpanName, ObservationType = 'agent') AS AgentNames,
groupArrayIf(LiteLLMRequestId, LiteLLMRequestId != '') AS RequestIds
FROM {database}.otel_traces
GROUP BY TeamId, ApiKeyHash, TraceId

View file

@ -6,7 +6,7 @@ use crate::Error;
const SCHEMA_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
const MIGRATIONS: [&str; 9] = [
const MIGRATIONS: [&str; 7] = [
include_str!("../migrations/0001_otel_traces.sql"),
include_str!("../migrations/0002_agent_traces.sql"),
include_str!("../migrations/0003_agent_traces_mv.sql"),
@ -14,8 +14,6 @@ const MIGRATIONS: [&str; 9] = [
include_str!("../migrations/0005_otel_traces_ttl.sql"),
include_str!("../migrations/0006_agent_traces_ttl.sql"),
include_str!("../migrations/0007_spend_logs_ttl.sql"),
include_str!("../migrations/0008_agent_traces_by_key.sql"),
include_str!("../migrations/0009_agent_traces_by_key_mv.sql"),
];
pub fn schema_statements(

View file

@ -40,8 +40,8 @@ async fn database() -> Result<Database, Box<dyn std::error::Error>> {
"CREATE DATABASE litellm",
"CREATE TABLE litellm.otel_traces (n UInt8) ENGINE = Memory",
"INSERT INTO litellm.otel_traces VALUES (1)",
"CREATE TABLE litellm.agent_traces (n UInt8) ENGINE = Memory",
"INSERT INTO litellm.agent_traces VALUES (2)",
"CREATE TABLE litellm.agent_traces_by_key (n UInt8) ENGINE = Memory",
"INSERT INTO litellm.agent_traces_by_key VALUES (4)",
"CREATE TABLE litellm.spend_logs (n UInt8) ENGINE = Memory",
"INSERT INTO litellm.spend_logs VALUES (3)",
"CREATE TABLE litellm.private_traces (n UInt8) ENGINE = Memory",
@ -86,6 +86,15 @@ async fn admin_sql_reads_rows_with_enforced_settings(
let json: Value = serde_json::from_str(&result)?;
assert_eq!(json["data"][0]["answer"], 1);
let result = read(
&database.client,
&connection,
"SELECT n AS answer FROM agent_traces_by_key",
)
.await?;
let json: Value = serde_json::from_str(&result)?;
assert_eq!(json["data"][0]["answer"], 4);
Ok(())
}

View file

@ -107,7 +107,7 @@ async fn schema_supports_span_rollups_and_spend_joins(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
let writer = Connection::writer(&database.url)?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
@ -144,7 +144,7 @@ async fn schema_supports_span_rollups_and_spend_joins(
let body = read_json(
&database,
"SELECT toUInt32(sum(SpanCount)) AS spans, toUInt32(sum(InputTokens)) AS tokens \
FROM trace_test.agent_traces WHERE TeamId = 'team-1' AND TraceId = 'trace-1'",
FROM trace_test.agent_traces_by_key WHERE TeamId = 'team-1' AND TraceId = 'trace-1'",
)
.await?;
assert_eq!(
@ -160,7 +160,7 @@ async fn keyed_rollup_keeps_same_trace_ids_separate_by_api_key(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
let writer = Connection::writer(&database.url)?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let timestamp = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64;
let rows = vec![
@ -204,7 +204,7 @@ async fn rollup_merges_spans_across_days_without_losing_root_fields(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
let writer = Connection::writer(&database.url)?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let day_start = time::OffsetDateTime::now_utc()
.replace_time(time::Time::MIDNIGHT)
@ -223,12 +223,16 @@ async fn rollup_merges_spans_across_days_without_losing_root_fields(
"ResourceAttributes": {"litellm.team_id": "team-1"}
}))?;
insert_rows(&database, "otel_traces", vec![child]).await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.agent_traces FINAL").await?;
execute_write(
&database,
"OPTIMIZE TABLE trace_test.agent_traces_by_key FINAL",
)
.await?;
let response = read_json(
&database,
"SELECT count() AS rows, any(RootName) AS RootName, any(RootInput) AS RootInput, \
any(RootStatus) AS RootStatus, sum(SpanCount) AS SpanCount \
FROM trace_test.agent_traces",
FROM trace_test.agent_traces_by_key",
)
.await?;
assert_eq!(
@ -247,7 +251,7 @@ async fn spend_deduplication_preserves_subsecond_requests_and_retries(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
let writer = Connection::writer(&database.url)?;
ensure_schema(&database.client, &writer, "trace_test", 7, 14).await?;
let now_ms = time::OffsetDateTime::now_utc().unix_timestamp_nanos() as i64 / 1_000_000;
let base_start_time = now_ms / 1000 * 1000;
@ -299,7 +303,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
#[future(awt)] database: TestResult<ClickHouseDatabase>,
) -> TestResult {
let database = database?;
let writer = Connection::writer(&database.url, "default", "")?;
let writer = Connection::writer(&database.url)?;
ensure_schema(&database.client, &writer, "trace_test", 30, 30).await?;
let old_time = time::OffsetDateTime::now_utc() - time::Duration::days(20);
let old_timestamp_ns = old_time.unix_timestamp_nanos() as i64;
@ -315,6 +319,7 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
}))?;
insert_rows(&database, "otel_traces", vec![span]).await?;
insert_rows(&database, "spend_logs", vec![spend]).await?;
assert_eq!(table_rows(&database, "agent_traces_by_key").await?, 1);
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
loop {
@ -337,10 +342,14 @@ async fn retention_changes_materialize_existing_rows_and_remain_idempotent(
tokio::time::sleep(Duration::from_millis(100)).await;
}
execute_write(&database, "OPTIMIZE TABLE trace_test.otel_traces FINAL").await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.agent_traces FINAL").await?;
execute_write(
&database,
"OPTIMIZE TABLE trace_test.agent_traces_by_key FINAL",
)
.await?;
execute_write(&database, "OPTIMIZE TABLE trace_test.spend_logs FINAL").await?;
assert_eq!(table_rows(&database, "otel_traces").await?, 0);
assert_eq!(table_rows(&database, "agent_traces").await?, 0);
assert_eq!(table_rows(&database, "agent_traces_by_key").await?, 0);
assert_eq!(table_rows(&database, "spend_logs").await?, 0);
let mutation_count = mutation_rows(&database).await?;
ensure_schema(&database.client, &writer, "trace_test", 14, 14).await?;
@ -359,7 +368,7 @@ async fn schema_statement_timeout_maps_to_transport_error() -> TestResult {
});
let client = Client::no_redirect_for_test();
let url = format!("http://{address}");
let writer = Connection::writer(&url, "default", "")?;
let writer = Connection::writer(&url)?;
let result = tokio::time::timeout(
Duration::from_secs(12),
ensure_schema(&client, &writer, "trace_test", 7, 14),

View file

@ -3,7 +3,6 @@ from typing import Final
from litellm.rust_bridge.traces import TraceStorage
OTEL_TRACES_TABLE: Final = "otel_traces"
AGENT_TRACES_TABLE: Final = "agent_traces"
AGENT_TRACES_BY_KEY_TABLE: Final = "agent_traces_by_key"
SPEND_LOGS_TABLE: Final = "spend_logs"