Skip to content

OpenObserve Explanation

What this page covers

How OpenObserve (O2) works inside and why it is built that way: node roles, deployment modes, the write and query paths, the durability model, storage formats and indexes, the DataFusion query engine, federated search, pipelines, the cost model behind the "140x" claim, and the security model. Exact defaults, ports, and prices are in Reference. Tasks are in How-to Guides.

Architecture Overview

OpenObserve is one Rust binary that can play several roles. In single-node mode one process does everything. In HA mode you run five main roles (Router, Ingester, Querier, Compactor, Scheduler) as separate Kubernetes workloads. They share object storage for data, PostgreSQL for metadata, and NATS for coordination.

The diagram shows the HA topology as the official Helm chart deploys it.

flowchart TB
    subgraph Sources["Telemetry sources"]
        OTEL["OTel Collector / SDKs<br/>(OTLP gRPC + HTTP)"]
        PROM["Prometheus<br/>(remote_write)"]
        ESC["Fluent Bit / Vector / Filebeat<br/>(_json, ES _bulk)"]
        RUMS["Browser and mobile RUM SDKs"]
    end

    subgraph O2["OpenObserve cluster (Helm chart)"]
        Router["Router<br/>(proxy + UI)"]
        Ingester["Ingester<br/>(WAL, Memtable, Parquet)"]
        Querier["Querier<br/>(DataFusion, caches)"]
        Compactor["Compactor<br/>(merge, retention, file list)"]
        Scheduler["Scheduler<br/>(alerts, reports, derived streams)"]
    end

    subgraph Deps["Shared services"]
        NATS["NATS<br/>(coordinator, events, queue)"]
        PG[("PostgreSQL<br/>(orgs, users, schemas, file list)")]
        OBJ[("Object storage<br/>S3 / GCS / Azure / MinIO / RustFS")]
    end

    Sources --> Router
    Router -->|ingest| Ingester
    Router -->|search| Querier
    Ingester -->|"Parquet / Vortex + .ttv index"| OBJ
    Querier -->|range reads| OBJ
    Querier -.->|"gRPC: unflushed data"| Ingester
    Compactor --> OBJ
    Scheduler -->|"scheduled SQL / PromQL"| Querier
    O2 --- NATS
    O2 --- PG

The main design choices:

  • Compute is separate from storage. All durable data lives in object storage as columnar files. Only ingesters hold state, and only for the short time before a flush. Queriers, compactors, routers, and schedulers can be replaced at any time.
  • Columnar files instead of a search index per document. Data is Apache Parquet (or, optionally, Vortex) compressed with Zstd. Full-text indexes, bloom filters, and partitions are opt-in per field. This is what separates OpenObserve from Elasticsearch, which indexes every field by default.
  • SQL first. Logs and traces use SQL (the Apache DataFusion dialect). Metrics use SQL or PromQL.
  • One binary for every signal. Logs, metrics, traces, RUM with session replay, continuous profiles (OTLP Profiles), synthetic checks, and LLM/agent traces share one storage and query layer.

Deployment Modes

OpenObserve runs in three shapes. The choice depends on durability and scale needs.

flowchart TD
    Start{"Need HA or<br/>horizontal scale?"}
    Start -->|No| Durable{"Need data to survive<br/>loss of the node disk?"}
    Durable -->|No| Local["Single node + SQLite + local disk<br/>(ZO_LOCAL_MODE=true, default)"]
    Durable -->|Yes| LocalS3["Single node + SQLite + object storage<br/>(ZO_LOCAL_MODE_STORAGE=s3)"]
    Start -->|Yes| HA["HA mode: Helm on Kubernetes<br/>PostgreSQL + NATS + object storage"]
    HA --> RBAC{"Need RBAC / SSO?"}
    RBAC -->|Yes| ENT["Enterprise image + OpenFGA + Dex<br/>(free up to 50 GB/day)"]
    RBAC -->|No| OSS["OSS image"]
Mode What it is good for Limits
Single node, local disk Dev, test, and light production. The docs measured about 31 MB/s (about 2.6 TB/day) on an Apple M2 with default settings No redundancy; the node's disk is the only copy
Single node, object storage Simple operations with durable, unbounded storage Still one process, so no HA for ingest or query
HA (cluster) Production at scale; every role scales on its own Needs Kubernetes, PostgreSQL, NATS, and object storage. Local disk storage is not supported

The vendor says the single binary scales to terabytes per day and HA mode to petabytes. The largest deployment it cites ingests more than 2 PB/day. These are vendor claims.

Node Roles

Each role is the same binary started with a different ZO_NODE_ROLE.

Role Responsibility State Scaling notes
Router Sends ingest requests to ingesters and search requests to queriers; serves the web UI Stateless Scale behind a load balancer
Ingester Parses records, runs ingest functions and real-time alerts, evolves schemas, writes WAL and Memtables, produces Parquet, and uploads it WAL, Memtable, Immutable, local Parquet Scale for write throughput; needs fast local disk (about 3,000 IOPS)
Querier Plans and runs SQL/PromQL with DataFusion over object storage plus ingester buffers Fully stateless; memory and disk caches only Scale for query concurrency and cache size
Compactor Merges small files into large ones, enforces retention, runs stream deletions, updates file-list indexes Stateless Scales horizontally; usually few nodes
Scheduler Runs scheduled alerts, reports, and derived streams; sends notifications Stateless Formerly named alertmanager (still accepted as an alias)

Workload isolation. ZO_NODE_ROLE_GROUP splits queriers into interactive and background groups. Alert, report, derived-stream, and search-job queries are classed as background, so a heavy scheduled report does not slow down a user's live search.

Write Path

The ingester keeps recent data in three places before it reaches object storage: the Memtable, Immutables, and local Parquet files that are not yet uploaded. Queriers read all three, so new data is searchable within seconds of ingest.

sequenceDiagram
    participant C as Client (OTLP / _json / _bulk)
    participant R as Router
    participant I as Ingester
    participant W as WAL (data/wal/logs)
    participant M as Memtable (Arrow)
    participant L as Local Parquet (data/wal/files)
    participant S as Object storage
    participant K as Compactor

    C->>R: HTTP or gRPC ingest request
    R->>I: Forward by stream
    I->>I: Parse, run ingest functions (VRL), set _timestamp
    I->>I: Evolve schema if needed, evaluate real-time alerts
    I->>W: Append records (hourly buckets)
    I->>M: Write Arrow RecordBatch
    Note over M: At ZO_MAX_FILE_SIZE_IN_MEMORY or WAL at ZO_MAX_FILE_SIZE_ON_DISK,<br/>Memtable becomes an Immutable
    M->>L: Every ZO_MEM_PERSIST_INTERVAL, dump Immutable as Parquet
    L->>S: Every ZO_FILE_PUSH_INTERVAL, merge a partition and upload<br/>when its size or age limit is reached
    S-->>K: Small files listed in the file list
    K->>S: Merge into files up to ZO_COMPACT_MAX_FILE_SIZE, delete old files
  • There is one Memtable (paired with one WAL file) per organization/stream_type.
  • Records get a _timestamp in microseconds. If a record has none, the ingester uses the current time. Events older than ZO_INGEST_ALLOWED_UPTO (5 hours by default) are rejected.
  • Schema changes (new fields, type changes) take a lock and update the stream schema in the metadata store. Past a field-count threshold, streams switch to a user-defined schema, where only listed fields are stored as columns.

Documented thresholds versus source defaults

The architecture page describes the thresholds as 256 MB in memory, 128 MB on disk, a 5 s persist interval, and a 10 s push interval. The v1.0.4 source defaults are 512 MB, 512 MB, 2 s, and 2 s. Upload also happens when a file reaches ZO_MAX_FILE_RETENTION_TIME (600 s). Check Reference and your own ZO_* settings rather than relying on either figure.

Durability Model

OpenObserve does not replicate in-flight data between ingesters. The docs argue that block storage and object storage are durable enough: EBS gp3 99.8%, EBS io2 (used by OpenObserve Cloud) 99.999%, S3 eleven nines. They also point out that cross-AZ replication on AWS costs about 2 cents per GB.

What this means in practice:

  • Before upload, data exists only on one ingester's volume (WAL plus local Parquet). If the process crashes, the WAL is replayed on restart. If the volume itself is lost, unflushed data is lost.
  • ZO_WAL_FSYNC_DISABLED defaults to true, so a node or kernel crash (not just a process crash) can lose the last WAL writes still in the page cache.
  • After upload, durability is that of the object store. Data is immutable. You cannot edit or delete individual records, only drop whole retention periods or streams.
  • RPO/RTO: because nodes are stateless or near-stateless, the vendor describes RPO and RTO as very low. RPO is bounded by the upload interval and file-age settings.

Correction

An earlier version of this page said a replication factor above 1 writes data to several ingesters. No such setting appears in the current docs or source. The documented design is a single in-flight copy.

Query Path

The querier that receives a search becomes the leader for that query. It fans work out to the other queriers over gRPC.

sequenceDiagram
    participant U as UI / API client
    participant R as Router
    participant QL as Leader querier
    participant PG as PostgreSQL (file list)
    participant QW as Worker queriers
    participant I as Ingesters
    participant S as Object storage

    U->>R: POST /api/{org}/_search (SQL, start_time, end_time)
    R->>QL: Route to a querier
    QL->>QL: Parse and validate SQL, derive time range
    QL->>PG: Look up files for stream and time partitions
    QL->>QL: Split the file list across queriers (consistent hash)
    par Fan-out over gRPC
        QL->>QW: Sub-plan with assigned files
        QW->>S: Range reads (cached on local disk)
    and Recent data
        QL->>I: Search Memtable, Immutable, local Parquet
    end
    QW-->>QL: Partial Arrow results
    I-->>QL: Partial Arrow results
    QL->>QL: Final merge, sort, limit
    QL-->>U: JSON result (or streamed partitions)

Query speed comes from reading less data:

  1. Time partitions. Files are laid out as org/stream/year/month/day/hour, so a narrow time range skips most files. Start and end time are required on every search.
  2. Custom partitions. KeyValue or hash partitions on low-cardinality fields (namespace, host) prune whole directories. They are immutable once set, and too many produce small files that hurt compression.
  3. Bloom filters skip files that cannot contain a high-cardinality value such as trace_id.
  4. Inverted index (Tantivy). match_all() uses a Tantivy index stored beside each data file (.ttv in an Iceberg Puffin container). The vendor reports up to 1000x faster full-text search for about 25% more storage.
  5. Column pruning and predicate pushdown. Parquet row-group statistics and projection let DataFusion read only the needed columns and row groups.
  6. Caching. Queriers cache files on local disk by default (up to 50% of free space, capped at 500 GB) and can also cache in memory. A result cache reuses overlapping earlier results. Ingesters can tell queriers to pre-cache newly uploaded files.

The vendor says partitioning, indexing, and caching reduce the search space "by up to 99% for most queries".

Query Engine: Apache DataFusion

Queriers run Apache DataFusion, the Rust query engine built on Apache Arrow. As of v1.0.4 OpenObserve pins DataFusion 54 from its own patched fork.

flowchart LR
    SQL["SQL or PromQL"] --> Parse["Parse and rewrite<br/>(match_all, histogram, UDFs)"]
    Parse --> LP["Logical plan"]
    LP --> Opt["Optimizer<br/>(predicate pushdown,<br/>projection and partition pruning)"]
    Opt --> PP["Distributed physical plan<br/>(leader + remote scan nodes)"]
    PP --> Scan["Parquet / Vortex scan<br/>+ Tantivy / bloom filter lookups"]
    Scan --> Exec["Vectorized execution<br/>(Arrow RecordBatches)"]
    Exec --> Merge["Leader merge<br/>(aggregate, sort, limit)"]
  • SQL extensions such as match_all, str_match, re_match, histogram, and approx_topk are DataFusion user-defined functions. The list is in Reference.
  • PromQL is implemented natively over the same metrics storage. v1.0 notes list "faster PromQL" as an improvement. Metrics can also be queried with SQL.
  • Streaming aggregations and partitioned queries let the UI show progressive results for long time ranges.
  • SIMD builds (latest-simd images) use AVX-512 or NEON for faster hashing and aggregation.

Storage Formats and Layout

Aspect Detail
Data files Apache Parquet by default. Vortex, a newer columnar format, can be selected per stream type, for example ZO_FILE_FORMAT=parquet,metrics=vortex
Compression Zstd for Parquet. Vortex uses OpenObserve's own UTF-8/Zstd compressor unless native compression is enabled
Layout files/<org>/<stream_type>/<stream>/<year>/<month>/<day>/<hour>/[partition keys]/<id>.parquet
Full-text index One .ttv Tantivy segment per data file
Metrics index During compaction of closed hours, a .midx file with sample blocks, labels, and a block directory (experimental, ZO_METRICS_INDEX_ENABLED)
Metadata PostgreSQL in cluster mode (SQLite in single-node mode): organizations, users, functions, alert rules, stream schemas, and the file list
Retention Per stream, enforced by the compactor; default ZO_COMPACT_DATA_RETENTION_DAYS=3650

Compaction matters because ingesters produce many small files, especially with many streams or partitions. The compactor merges them into large, time-sorted files (up to 2 GB by default) and records the new files in the file list. Each stream compacts on its own, so thousands of streams mean more concurrent merge jobs and more file-list rows in PostgreSQL. No public benchmark covers 10,000+ streams.

Multi-Tenancy

Organizations are the tenant boundary. Each has its own streams, dashboards, alerts, pipelines, functions, and members, and every API path starts with /api/{org}/. A user can belong to several organizations.

Recent releases extended this:

  • v0.91.0 added a Super Org model, org-level ingestion tokens, and org-level storage configuration, so different tenants can write to different buckets.
  • Enterprise QoS (query and workload management) limits and prioritizes query resources per tenant.

Federated Search (Super Cluster)

Enterprise Super Cluster joins several OpenObserve clusters, typically one per region, each with its own bucket. It works like the in-cluster query path, one level up:

  1. The cluster that receives a search becomes the leader cluster for it.
  2. It finds the other clusters from super cluster metadata and calls each over gRPC with the same query.
  3. Each worker cluster runs the normal leader/worker querier flow locally.
  4. The leader cluster merges the results and returns them.

Configuration metadata such as stream schemas, dashboards, and alerts is synchronized across member clusters (over NATS), while ingested data stays in the cluster that received it, which helps with data residency (federated search architecture, checked 2026-09-28).

Pipelines and Functions

Pipelines transform data at ingest or on a schedule. They are built in a visual editor from source, transform, and destination nodes:

  • Transforms are VRL (Vector Remap Language) functions and conditions. Uses include parsing, enrichment from enrichment tables, PII redaction, dropping noise, and normalizing fields.
  • Destinations can be other streams, which allows routing and logs-to-metrics conversion, or external endpoints.
  • Scheduled pipelines and derived streams run SQL on a schedule and write the results to a new stream. They run on the Scheduler.

VRL at ingest costs CPU on the ingesters. The performance docs say to test throughput with your real functions.

Product Surfaces Beyond Logs, Metrics, and Traces

Surface What it does
RUM Browser (and mobile) SDKs for Core Web Vitals, errors, performance, and session replay
Continuous profiling CPU, memory, and lock profiles over OTLP Profiles (async-profiler, pprof, OTel eBPF profiler)
Service graph Service-to-service dependencies with per-edge request counts and health colors
Alerts and incidents Scheduled, real-time, anomaly, composite, and SLO burn-rate alerts; related alerts grouped into incidents with a lifecycle
Synthetic monitoring Browser and HTTP checks with private locations (v0.92.0); moved to open source in v1.0
AI observability LLM and agent tracing, cost and token tracking, evaluations with LLM-as-judge, annotation queues, and an agent graph (v0.92.0 to v1.0)
O2 AI Assistant In-product assistant that writes SQL, VRL, and PromQL

The "140x Cheaper" Claim

The README says storage costs up to 140x less than Elasticsearch. The claim combines several effects. The breakdown below is this vault's model, not the vendor's published method.

Factor OpenObserve Elasticsearch (self-managed) Rough effect
Compression Columnar Parquet + Zstd, no index on every field Lucene segments plus inverted indexes on all fields Several times fewer bytes stored
Storage tier Object storage (S3 Standard about $0.023/GB-month) Block SSD (roughly $0.08-0.10+/GB-month) About 4x cheaper per GB
Replication None in the application; the object store is durable Usually one or more replica shards About 2-3x less
Idle compute Stateless queriers can scale down Data nodes stay up to hold shards Compute, not storage

Multiplying 5-10x (compression) by about 4x (tier) by 2-3x (replicas) gives roughly 40-120x, the same order as the vendor figure. How much you save depends on data entropy, which fields you index, and object storage request costs.

Caveats

  • The 140x figure and "a quarter of the hardware" are vendor claims. No independent benchmark was found.
  • Object storage request charges (GET/PUT) and egress can matter for query-heavy workloads.
  • Full-text search without an inverted index is a scan. Enabling Tantivy on text fields adds about 25% storage.
  • Earlier versions of this page gave an illustrative 100 GB/day, 30-day comparison: about $507/month for OpenObserve versus about $1,800/month for Elasticsearch. It was not sourced. Treat it as a rough model only.

Comparison with Elasticsearch

Aspect OpenObserve Elasticsearch
Storage Parquet/Vortex on object storage Lucene segments on local disks
Indexing Opt-in per field (bloom filter, Tantivy inverted index, partitions) Inverted index on every field by default
Aggregations Columnar and vectorized (Arrow), strong on analytics Doc values; fast, but uses more memory
Wildcard and fuzzy text search Fast only with an inverted index on the field; otherwise a scan Native strength
Query language SQL, PromQL Query DSL, ES|QL, KQL/Lucene
Ingest API compatibility Accepts _bulk from Beats, Vector, and Fluent Bit Native
Mutability Immutable records; retention drops only Update and delete by query

Security Model

Authentication

  • Every API call uses HTTP Basic auth. Users send email:password; service accounts send email:token. OTLP/gRPC sends the same header plus organization and stream-name metadata.
  • The root user is created from ZO_ROOT_USER_EMAIL and ZO_ROOT_USER_PASSWORD on first start. There is one per installation, with access to all organizations. ZO_ROOT_USER_TOKEN optionally sets a root API token.
  • SSO (Enterprise and Cloud) goes through Dex, an OIDC/OAuth2 broker, to LDAP, OIDC, SAML, and social providers. Dex group and role attributes map users to organizations and roles (O2_DEX_GROUP_ATTRIBUTE, O2_DEX_ROLE_ATTRIBUTE).

Authorization

  • Open source: no RBAC. Organization roles are root, admin, and member, and members can access all data in their organizations. Service accounts get full access.
  • Enterprise: RBAC is backed by OpenFGA. OpenObserve stores role relationships in OpenFGA and asks it a true/false question for each action. There are four predefined roles (Admin, Editor, Viewer, User), custom roles with per-resource permissions (list, get, create, update, delete), user groups, and service accounts with no permissions until a role is assigned. RBAC requires HA mode.
  • Cloud: RBAC is preconfigured. Service accounts are not supported.

The request flow below shows how Enterprise checks authentication and authorization.

sequenceDiagram
    participant B as Browser / collector
    participant R as Router
    participant D as Dex (OIDC broker)
    participant IdP as IdP (LDAP / OIDC / SAML)
    participant N as Querier or Ingester
    participant F as OpenFGA
    participant PG as PostgreSQL

    alt Browser SSO login
        B->>R: Open UI
        R->>D: Redirect to Dex
        D->>IdP: Authenticate user
        IdP-->>D: Identity + groups
        D-->>R: ID token (groups, role attributes)
        R->>PG: Map groups to org and role
    else API or collector
        B->>R: Authorization: Basic email:token
        R->>PG: Validate user or service account
    end
    R->>N: Forward request with user context
    N->>F: Check(user, action, resource)
    F-->>N: allowed = true / false
    N-->>B: Result or 403

Threat Model Notes

Threat Built-in mitigation Gap or operator action
Credential theft Service account tokens are shown once and can be rotated Basic auth sends credentials on every request, so always use TLS
Network sniffing Built-in TLS for HTTP and gRPC (ZO_HTTP_TLS_*, ZO_GRPC_TLS_*), off by default Enable TLS or terminate it at the ingress
Rogue node joining the cluster ZO_INTERNAL_GRPC_TOKEN for internal gRPC Keep 5081 internal; use NetworkPolicies
Bucket exposure None in the app; data sits as Parquet in the bucket Encrypt the bucket (SSE-KMS), restrict with bucket policies, use IRSA
Cross-tenant access Organization scoping on every API; RBAC in Enterprise The open-source edition has no per-stream access control within an org
PII in logs Enterprise Sensitive Data Redaction; VRL redaction in pipelines (all editions) Redact at ingest; data is immutable once written
Metadata DB compromise Store DSN in a Secret (the chart supports this) sslmode=require, least-privilege DB user

Design Trade-offs

Decision Benefit Cost
Object storage for all data Cheap, durable, unbounded storage; stateless compute Object-store latency on cold reads; request costs; needs caches
No in-app replication Simpler, and no cross-AZ transfer fees Unflushed data depends on one ingester volume
Opt-in indexing High compression, fast ingest (the vendor says 5-10x faster than Elasticsearch on the same hardware) You must choose partitions, bloom filters, and full-text fields per stream
Immutable data Integrity for audit and compliance No record-level delete or update; GDPR erasure is by retention or stream
Single binary, many roles Easy local start; the same code in HA Many features in one release train; fast release pace
AGPL-3.0 core plus commercial Enterprise Keeps changes to the core open SaaS embedding needs a legal review or a commercial license; RBAC and SSO are paid above 50 GB/day

Sources