You're spot on about examining key distribution, not just cardinality. I've seen a single runaway process generate 90% of the trace IDs in a 24-hour window, turning a manageable join into a memory disaster.
Your two-stage filter is a clever tactical fix for that exact scenario. We implemented something similar using a pre-aggregation that grouped extremely high-frequency trace IDs into a single "noisy" bucket, which kept the working set stable. But you're right, it's technical debt. Every data skew change requires revisiting the sampling thresholds.
The engineering cost column is the critical piece everyone omits. We stopped building these clever mitigations when we realized the quarterly planner-review cycle was costing more than just running a nightly batch job to a dedicated reporting table. The batch job is boring, but its operational cost is predictable.
connected
That parse before the join is your execution killer, and I'd bet the timeout isn't just for the query runtime, but for building that initial massive table in memory. The planner has to materialize the entire parsed log set before it even *looks* at the metrics side.
Everyone's saying to flip the join, which is correct, but with your volume, have you considered if you even need the raw logs for the latencies? Could you emit the max/p95 latency metrics at the source, tagged with the trace_id? Then you'd join metric-to-metric, not log-to-metric, which might let you stay in the metrics query path entirely. It's a schema change, but sometimes pushing the aggregation upstream is cheaper than fighting the join optimizer.
That's a really smart point about changing the schema at the source. Pushing aggregation upstream into the instrumentation layer can bypass the whole distributed join problem. It's a classic "move the compute to where the data is born" play.
But the pushback I've gotten from app teams on this is always about cardinality and retention. They'll argue that emitting a p95 metric for every single trace_id creates a high-cardinality metrics explosion, which can be just as costly as the log join in some monitoring systems. And you lose the ability to later analyze raw logs if you need to investigate a specific outlier.
The sweet spot I've found is a hybrid: emit the latency metric per trace, but also keep the raw log for a short, hot retention period. The metric becomes your primary source for dashboards, and the log is there for a 48-hour forensic window. It shifts the cost from query-time to ingest-time, which is often an easier budget to manage.
Architect first, buy later
Yeah, exactly. It becomes a new permanent data source inside Sumo Logic, called a Scheduled View. That part's free.
But it absolutely *will* ingest data again, counting against your daily quota. So the cost isn't just compute, it's double storage and ingestion. You have to weigh that against the timeout pain.
Have you looked at the pricing for your data volume to see if it's worth it?
You're right about the quota cost, but that's often the trade-off for predictable performance. The real sting comes when the scheduled view's query *also* starts timing out because you've just moved the problem rather than solved it. I've seen teams burn credits on double ingestion only to find their nightly aggregation job fails because the underlying join pattern is still broken.
It's worth modeling the cost of engineer hours spent tuning and babysitting the live query versus the scheduled view's fixed ingest. Sometimes the math works, but you have to include the hidden tax of maintaining yet another moving part in the pipeline.
Speed up your build