# Streama Explorer

An interactive tour of Streama, the Kafka-based engine that parses, routes, alerts on and stores Coralogix telemetry as it streams.

Source: https://benchmarks.coralogix.com/streama

## Welcome to the Streama Tour

Streama is Coralogix's Kafka-based processing engine. It parses, enriches, transforms, alerts on and generates metrics from data **as it streams**, before anything is indexed or stored, in three phases. **Source** does all the input I/O up front and lands events on Kafka. **Stream** works only on data already in memory, so it is CPU-bound and scales on compute alone. **Sink** does the output I/O at the end: your cloud object storage, and everything that reads it.

## In-stream means CPU bound

Stream processing is **CPU-bound** work, so Coralogix scales ingestion and processing independently of storage and indexing. When traffic spikes, Streama adds capacity in seconds and releases it again as traffic falls. **This makes Coralogix highly efficient to run, and allows us to pass those savings on to our customers.**

## Data arrives, and lands on Kafka

Logs, metrics, traces, profiles, RUM sessions and security events arrive raw from your apps, clusters and clouds, and from the **code agents** writing your software. The **ingest gateway** authenticates each event, and Kafka Connect produces it to a **Kafka** topic: partitioned and replicated, so many workers read in parallel, order holds within each partition, and a node failing loses nothing. **Durable and persistent**, and the last I/O before processing.

## Raw lines become structured fields

Parsing rules turn each raw line into a **structured record**: they parse JSON and key-value text, extract fields with regular expressions, set the event's timestamp, and replace or **mask sensitive values** before anything downstream sees them. Every event leaves Parse with named, typed fields, ready for every step after it.

## Not a static index. An evolving schema

Most platforms fix a field's type the first time they see it. When the type changes, as it does whenever a team ships new code, you get a **mapping exception**, and the event is rejected or the field goes unindexed. The **schema registry** tracks every field's type over time instead and supports **multi-type fields**, so a field can be a number one day, a string the next and an object after that. **Every event is kept**, and every version of the field stays queryable.

## The TCO Optimizer decides first

Straight after parsing, policies match on application, subsystem, severity or a DPXL expression, evaluated top to bottom: **the first match wins** and sets a priority. **High** gets everything and lands in OpenSearch, **Medium** gets monitoring and lands in your bucket, and **Low** is kept for compliance. **Blocked** is dropped here and costs nothing. Metrics have their own optimizer, **Dynamic Metrics TCO**, which routes each series by usage: **Active** goes on for full analysis, while **Historical** is written to your bucket for retention and query.

## Then everything happens at once

After **Enrich** adds context, security intelligence, AWS cloud metadata, Kubernetes metadata, geo-IP and custom enrichment that adds your own business context, each event fans out across parallel lanes in the same instant. **Alerting** analyses data as it moves through the stream: instant triggers, no indexing lag, and alerts keep operating even when queries or ingest lag. Beside it, **Loggregation** clusters millions of lines into a few hundred templates, **Events2Metrics** generates a metric from a log and can drop the original, the **metrics usage analyser** blocks series nobody reads before **recording rules** derive lean ones, and the **Multi-SLM evaluation engine** scores every AI prompt and completion. All of it runs before storage, so nothing waits on indexing, and this is only a subset of what Coralogix can do.

## AI spans, classified in flight

Spans from your AI applications carry each prompt and completion. As they stream, the AI Center's **Multi-SLM architecture** evaluates them: a bench of compact, transformer-based **small language models**, each purpose-trained for one judgement. **Prompt injection**, **PII**, hallucination (**context adherence**, **correctness**, **completeness**), **toxicity** and **allowed or restricted topics** are a few, and your own custom evaluators sit beside them. Every enabled evaluator scores every span, without adding latency to ingestion or touching the live request. The scores become **labels on the span**: high scores surface as issues on the application's dashboard, and every score is queryable in DataPrime, so a flagged interaction can raise an alert.

## Your storage is not an archive

Streama writes every event to **your own cloud object storage** by default, as **wide Parquet**, a columnar format tuned for observability: one column per field, grouped into row groups. DataPrime queries it there directly, the moment it lands, with no rehydration and no re-indexing. You route data into your own **datasets** and choose the bucket each one writes to. Each dataset carries its own schema and permissions, so you get **multi-tenancy inside your own data lake**. Events land as Parquet and metrics in an open format, so any engine that reads Parquet reads your telemetry, Apache Spark, Trino, Athena, DuckDB or your warehouse: **a nexus, not a dead end**. For the small subset of events that need millisecond ingestion and the fastest queries, you can also write to **OpenSearch**, and chosen copies are forwarded over Kafka to your SIEM or data lake.

## Data is never at rest

Dynamic materialization weighs each field's cardinality and how often it is read. **The fields you query most stay materialized** for speed, and the ones you rarely touch move into a lower-footprint encoding, still fully queryable: DataPrime decodes them on the fly. A query reads **only the columns it needs**, straight from the bucket. Even while a column is materialized, Coralogix still delivers **5x compression** for events and **30x** for metrics. Watch as some columns materialize while others are compressed.

## The DataPrime query engine

A highly parallelised query engine that plans every query and runs it across many nodes, reading **Parquet straight from your bucket** with no rehydration and no re-indexing. One language for all your telemetry: filters and extracts, **joins** across logs and spans, **unions** across datasets, groupby aggregations and **window functions**, designed to scan huge, high-cardinality data efficiently and quickly. It serves Explore, dashboards, alerts, Olly and your tools, and what an investigation finds can be written back into the lake as a new **dataset**.

## The distributed metrics engine

Fully **Prometheus-compatible**, and a highly parallelised metrics engine written in **Rust**, designed to plan, execute and read **huge-cardinality** queries efficiently and quickly. It serves **PromQL** to dashboards, SLOs, alerts, Olly and your tools: rates, aggregations, quantiles and downsampling, spread across many nodes.

## For agents and engineers

In front of your telemetry lake, our interfaces offer unlimited access to your data. Whether through our **products**, such as APM, RUM, SIEM, the AI Center and more, or **Olly**, our autonomous observability agent, which combines a knowledge base, a reasoning model, skills, memories and rules with a bench of specialist agents. Or bring your own coding agent, connecting through our **MCP server** or the **CLI**, for low-token, high-impact exploration of your telemetry lake.

## Now explore it yourself

Rotate and zoom the model, hold on any stage to fly to it (or double-click), and click any stage for detail.
