Exactly. The maintenance hazard is a form of technical debt that accrues interest silently. That hidden exclusion list isn't just a query parameter, it's a business rule disguised as a performance fix. Its correctness depends entirely on the correlation between your performance filter (e.g., high cardinality trace IDs) and the business logic you're analyzing.
A concrete example: excluding high-cardinality `trace_id` values to speed up a latency join assumes those traces are uniformly distributed noise. But if a new service deployment suddenly creates high-cardinality traces that are also high-latency outliers, your filtered analysis becomes misleading precisely when you need it most. The verification queries to catch this drift become as complex as the original problem.
The scheduled view forces this cost into the design phase, where you explicitly model the data reduction. It's more expensive upfront, but its failure modes are contained within the pipeline's own SLA, not hidden inside an analytical query's semantics.
brianh
Yeah, the `receiptTimeout` increase is a classic trap. You're just moving the goalposts on a query that's fundamentally trying to process too much data in a single pass. The in-memory threshold isn't undocumented by accident; it's a hard system limit to prevent runaway queries from taking down the cluster.
Your real enemy is that initial parse operation on the entire 72-hour log window. You're extracting a high-cardinality key from every single log line before you even think about joining. A more surgical approach is to push the filtering and aggregation logic into a subsearch, forcing a reduction before the join. For example, pre-aggregate your metrics side by `trace_id` first, or use a lookup to map trace IDs to services you care about before the heavy join. It forces the planner to work with a smaller, summarized dataset.
Have you explored using the `lookup` operator with a scheduled lookup table as a join alternative? It can sometimes bypass the in-memory join issue entirely by acting as a distributed hash table, though the setup is more involved.
api first
Pushing the parse and filtering into a subsearch is the right instinct, but I've found it often doesn't force the planner's hand like you'd hope. The optimizer can still decide to materialize the full, pre-parsed dataset before applying your filters if the cost model says so.
The lookup table idea can work, but you're trading one set of problems for another. You now have to manage the freshness and cardinality of that lookup table. If your trace_id space is huge and volatile, the lookup itself becomes a large, hot dataset that needs to be broadcast, reintroducing the memory pressure you were trying to avoid.
Sometimes the only fix is to bite the bullet and materialize an aggregated fact table on a schedule. It's more upfront work than a clever subquery, but it's predictable.
-- bb
Exactly. The optimizer's cost model is a black box and it loves to "help" by undoing your careful subquery. I've seen it decide to fully materialize a 50GB intermediate dataset because its outdated stats said the filter selectivity was low, when the real data would have cut it down to 5GB.
That aggregated fact table is the only real fix. All the clever tricks just become a maintenance nightmare when the planner changes its mind after an engine upgrade.
Trust but verify.
The cost of a planner surprise after a minor engine patch can dwarf the scheduled aggregation job. It's not just about correctness, it's about operational stability.
Your 50GB to 5GB example is classic. I've had to add hard query hints as a stopgap, which is another form of technical debt that breaks on upgrade. The aggregated table is the only way to lock in the predictable cost.
Show me the bill
Hilarious that you think they'd grant a custom capacity increase. They'd just suggest the scheduled search you're already describing, and charge you more for it.
That move out of the interactive layer is the real answer, but "pre-aggregate your metrics into a timeslice before the join" is optimistic. If your join key is high cardinality, your timesliced aggregation is still huge. You're just moving the memory wall, not removing it.
The batch job to a separate index is the only way. Everything else is hoping the black box behaves.
Keep it simple
That parse operation on a 72-hour log window is the killer. You're forcing the engine to scan and extract `trace_id` from every single log line before it can even think about joining, which creates a massive intermediate dataset right out of the gate.
I hit this same wall with Datadog Logs. The fix for us was to invert the logic: pre-aggregate the metrics side first into a much smaller dataset, *then* join. Can you push the `timeslice` and aggregation into the metrics subquery? Something like:
```sql
dataset=metrics _metric=app_latency_ms
| timeslice 1h
| max(_value) as max_latency, pct(_value, 95) as p95_latency by _timeslice, trace_id
| fields _timeslice, trace_id, max_latency, p95_latency
```
Then join your logs to that. It reduces the join from a massive many-to-many to a (hopefully) smaller many-to-one. If the metrics dataset is still too big, you might need that aggregated fact table everyone's mentioning.
Dashboards or it didn't happen.
Increasing the receiptTimeout is just hitting snooze on a broken query. That parse on a 72-hour window is forcing a full scan and extract before anything else, creating a massive intermediate set that's destined to fail.
The real fix is to never parse trace_id from all logs. Pre-aggregate the metrics side into a tiny dataset first, then join. Your example still joins before the timeslice, which is the core mistake. Invert it: aggregate metrics by trace_id and *then* join, not the other way around.
That's a good point about isolating the metrics cardinality first. I'm actually dealing with something similar right now, where a 24-hour window is blowing up our cluster. I ran the isolation query and it's coming back at 12 million distinct trace IDs, which explains everything.
But about pre-filtering the metrics dataset by time before aggregation - doesn't that defeat the purpose of a scheduled view? If I'm building a view for a dashboard that looks at the last 30 days, I'd have to schedule 30 separate aggregations, one for each day, right? Or am I misunderstanding how the pre-aggregation would work?
Learning by breaking
That 12 million distinct trace ID count clarifies a lot. On your pre-filtering question, you wouldn't schedule 30 separate aggregations. The scheduled search should build a single, cumulative aggregated table.
You'd run a batch job, say daily, that processes only the *new* metrics data from the last 24 hours, aggregates it by trace_id and date, and appends those results to a persistent summary index. Your dashboard query then runs against that pre-aggregated index, scanning dramatically fewer records. The time filter for "last 30 days" is applied to this small summary table, not the raw data. Does your platform support writing the results of a scheduled search to a permanent index?
Yep, that's the classic timeout spiral. Increasing `receiptTimeout` is a band-aid because the `parse` on that massive log set is the real killer. It has to scan every raw log line before the join even starts, building a huge table in memory that's doomed.
Your posted query pattern is backwards. The metrics dataset is the smaller one, right? So you need to aggregate *that* side first, down to maybe a few hundred thousand rows, *then* join it to your logs. Run the metrics query as a scheduled search, dump the results to a lookup, and join to that. It flips the memory problem on its head.
But honestly, with 12 million distinct trace IDs like someone else mentioned, even a pre-aggregated metrics lookup might be too big for a live join. That's when you gotta move to a batch process and a separate reporting index. The interactive query layer just can't handle that scale without exploding.
You're right that flipping the join order is the only sane move, but calling it a "band-aid" is a bit much. It's more like tourniquet-level surgery that still might not save the limb.
The real kicker is when that aggregated metrics "lookup" table itself hits 12 million rows, which your post hints at. At that point, the planner can still decide to do something spectacularly dumb, like hash-joining it in memory against the full log scan anyway. I've watched a vendor's query engine do exactly that, merrily blowing through a terabyte of cluster memory because its stats said the lookup was "small". The batch process and separate index aren't just a scaling step, they're a jailbreak from the optimizer's whims.
cg
Exactly. I stopped trusting planner stats years ago. You can have perfect indexes and a tiny lookup table, and the engine will still decide to broadcast it to every node for a nested loop join against your terabyte fact table.
That's why my "band-aid" comment stands. Flipping the join is surgery, sure, but you're still operating inside the unpredictable black box. The separate index isn't just scaling, it's removing the planner from the equation entirely. You query the pre-joined table. End of story.
SQL is enough
Increasing the receiptTimeout is a classic move, but it's like hoping the engine can swim after you've already thrown it in the deep end. The real issue is your query's fundamental shape. That parse operation is forcing a full scan of a multi-terabyte log set before any join happens, creating a massive intermediate table that's dead on arrival.
Everyone's correctly pointing out to flip the join order, but they're missing the follow-through. Even if you pre-aggregate the metrics side first, you're still at the mercy of Sumo's optimizer when you join back. I've seen it take a perfectly reasonable 500k-row lookup and decide to broadcast it across the entire log scan, recreating the memory wall you tried to avoid. The scheduled search to a separate index isn't just a performance tweak, it's the only way to remove the planner from the equation entirely. Are you able to write scheduled search results to a permanent index, or are you stuck in the interactive layer?
Good call on checking the plan with `explain`. It's the first thing I do when a join acts up.
But in my experience, hints like `/*+ BROADCAST(metrics) */` can be a double-edged sword. If the underlying data volume shifts and the "small" table grows, that hardcoded hint can blow up spectacularly later. It locks you into a manual review cycle.
I've started logging execution plans from our dashboards to catch when a hinted query outlives its assumptions. It's a bit of work, but it beats a 3am page because a broadcast table got too big.
data over opinions