Chapter 10 · Data Streams, Time-Series Data, Logs, Metrics, and Append-Heavy Workloads
Design a Telemetry Data Stream with Lifecycle, Schema, Routing, Search, Aggregation, and Retention Requirements
Assemble a production-style AtlasMart telemetry design that joins schema, rollover, lifecycle, routing, querying, aggregation, security, restore requirements, observability, and acceptance tests into one deployable contract.
Learning outcomes
The chapter ends with a design exercise: AtlasMart needs one telemetry contract that engineers can deploy, observe and review. The goal is not a magical “best” shard size or retention period. It is a specification that connects schema and lifecycle decisions to measurable workload requirements and makes Elastic/OpenSearch divergences explicit.
Write a telemetry-stream manifest that assigns ownership and pins the platform versions/portability boundary.
Define schema, rollover, late-data, correction, retention, security and restore acceptance criteria together.
Build deterministic ingest/search/aggregation tests that span a rollover boundary.
Measure indexing lag, storage growth, shard count and query latency rather than using universal tuning folklore.
Produce separate lifecycle adapters for Elastic DLM/ILM and OpenSearch ISM without disguising them as one API.
The reproducible examples target self-managed Elasticsearch 9.5.3 and OpenSearch 3.8.0 using the established AtlasMart lab conventions: Elasticsearch on https://localhost:9200 with the copied CA certificate, OpenSearch on https://localhost:9201 with the disposable demo certificate explicitly treated as local-only, one-node disposable clusters, pinned versions, and no moving latest tags. The final architecture keeps ordinary data streams as the shared baseline. Elastic data stream lifecycle/ILM and OpenSearch ISM are lifecycle systems with different APIs and semantics; the application contract should describe retention intent, while platform adapters implement it per product.
The generation environment does not run the two search servers, so commands below are deterministic lab specifications and expected invariants, not fabricated captured output. Run them only against disposable AtlasMart course resources. Record your actual responses, timestamps, backing-index names, latency, shard counts, and disk usage before drawing operational conclusions.
1. Start with an explicit design manifest
telemetry_stream: atlasmart-telemetry
owners:
schema: search-platform
ingest: observability-platform
producer: atlasmart-services
servers:
elasticsearch: 9.5.3
opensearch: 3.8.0
portable_contract:
timestamp: "@timestamp"
shards: 1 # lab only
replicas: 0 # lab only
schema: strict
write_model: append-heavy
rollover: policy-driven in production
correction: privileged backing-index workflow
quality:
required_fields: ["@timestamp", event_id, event_kind, service, environment, tenant_id]
reject_unknown_fields: true
tenant_scope_required: true
slo_inputs:
ingest_rate: measured
indexing_lag_p95: measured
search_latency_p95_p99: measured
storage_bytes_per_day: measured
late_event_rate: measured
retention:
live_search_window: business-defined
snapshot_window: recovery-defined
legal_hold: separately governed
The 1 shard / 0 replicas values are course-lab
assumptions, not production recommendations. Production
shard/replica counts depend on ingest rate, availability,
storage, recovery and query concurrency. The manifest therefore
separates lab reproducibility from measured production inputs.
2. Deploy the portable schema and stream
PUT _index_template/atlasmart-telemetry-template-v1
{
"index_patterns": ["atlasmart-telemetry"],
"priority": 300,
"data_stream": {},
"template": {
"settings": {
"number_of_shards": 1,
"number_of_replicas": 0
},
"mappings": {
"dynamic": "strict",
"properties": {
"@timestamp": {"type":"date"},
"event_id": {"type":"keyword"},
"event_kind": {"type":"keyword"},
"service": {"type":"keyword"},
"environment": {"type":"keyword"},
"tenant_id": {"type":"keyword"},
"host_id": {"type":"keyword"},
"level": {"type":"keyword"},
"message": {"type":"text"},
"duration_ms": {"type":"double"},
"request_count": {"type":"long"},
"error_count": {"type":"long"},
"labels": {"type":"object", "dynamic":"strict"}
}
}
}
}
PUT _data_stream/atlasmart-telemetry
Before ingest, assert the resolved mapping and template version. Chapter 09 established simulation/template discipline; carry it forward because a data stream creates future backing indices from templates. A template drift today can appear only after tomorrow's rollover.
POST _index_template/_simulate_index/atlasmart-telemetry
GET _index_template/atlasmart-telemetry-template-v1
GET _data_stream/atlasmart-telemetry
3. Acceptance fixture spanning rollover
POST atlasmart-telemetry/_doc?refresh=true
{"@timestamp":"2026-09-11T07:00:00Z","event_id":"acc-1","event_kind":"log","service":"catalog-api","environment":"lab","tenant_id":"tenant-a","host_id":"node-1","level":"INFO","message":"before rollover","duration_ms":10.0,"request_count":1,"error_count":0,"labels":{}}
POST atlasmart-telemetry/_doc?refresh=true
{"@timestamp":"2026-09-11T07:01:00Z","event_id":"acc-2","event_kind":"metric","service":"catalog-api","environment":"lab","tenant_id":"tenant-a","host_id":"node-1","level":"INFO","message":"minute metric","duration_ms":11.0,"request_count":100,"error_count":1,"labels":{}}
POST atlasmart-telemetry/_rollover
POST atlasmart-telemetry/_doc?refresh=true
{"@timestamp":"2026-09-11T07:02:00Z","event_id":"acc-3","event_kind":"log","service":"checkout-api","environment":"lab","tenant_id":"tenant-b","host_id":"node-2","level":"ERROR","message":"after rollover","duration_ms":90.0,"request_count":1,"error_count":1,"labels":{}}
Then prove the logical stream returns all three documents while
the physical _index field shows at least two
generations.
GET atlasmart-telemetry/_search
{
"track_total_hits":true,
"sort":[{"@timestamp":"asc"}],
"query":{"ids":{"values":[]}},
"aggs":{}
}
# Replace the empty ids query with a terms query on event_id if desired:
GET atlasmart-telemetry/_search
{
"track_total_hits":true,
"query":{"terms":{"event_id":["acc-1","acc-2","acc-3"]}},
"sort":[{"@timestamp":"asc"}]
}
The first skeleton intentionally contains an empty IDs clause as
a review point; use the second concrete request for the
acceptance test. The invariant is exact count 3 and timestamps
in expected order, not fixed _score values.
4. Facet and time-series acceptance query
GET atlasmart-telemetry/_search
{
"size":0,
"query":{"range":{"@timestamp":{"gte":"2026-09-11T06:59:00Z","lt":"2026-09-11T07:03:00Z"}}},
"aggs":{
"by_service":{"terms":{"field":"service","size":10}},
"events_over_time":{"date_histogram":{"field":"@timestamp","fixed_interval":"1m","min_doc_count":0}},
"p95_duration":{"percentiles":{"field":"duration_ms","percents":[95]}},
"total_requests":{"sum":{"field":"request_count"}},
"total_errors":{"sum":{"field":"error_count"}}
}
}
On this tiny fixture you can assert exact bucket counts and sums. Do not turn the percentile number into a production latency benchmark; percentiles and distributed aggregation behavior need representative scale.
5. Lifecycle adapters: shared intent, different implementation
| Intent | Elasticsearch 9.5.3 | OpenSearch 3.8.0 |
|---|---|---|
| Automatic rollover/retention | Data stream lifecycle and/or ILM depending on design | ISM policies applied to backing indices/data stream workflow |
| Data-stream-level retention API | Built-in data stream lifecycle supports data_retention | No assumption of Elastic DLM API; model with ISM/target platform controls |
| Time-series-specialized mode | TSDS with index.mode=time_series and dimension/metric mapping | Do not assume equivalent syntax |
| Timestamp field | @timestamp for ordinary Elastic data streams | Defaults to @timestamp; OpenSearch template can specify a custom timestamp field |
| Modify backing membership | Elastic data stream modify APIs | OpenSearch 3.8 adds experimental modify-data-stream API |
Keep the lifecycle intent in source control: maximum searchable age, restore obligations, rollover SLO and cost boundary. Render it to each platform's supported lifecycle configuration instead of maintaining “equivalent-looking” JSON that is semantically different.
6. Measure before choosing rollover thresholds
GET _data_stream/atlasmart-telemetry/_stats
GET _cat/indices/.ds-atlasmart-telemetry*?h=index,docs.count,pri.store.size,store.size
GET .ds-atlasmart-telemetry*/_stats/indexing,search,store,segments
GET _cat/shards/.ds-atlasmart-telemetry*?v
Capture ingest documents/s, bytes/s, indexing lag, search p95/p99, shard store growth, segment counts, recovery duration and snapshot duration under representative load. Rollover thresholds balance shard size, shard count, recovery time, merge work and search fan-out; no universal threshold replaces those measurements.
7. Failure and restore tests
| Test | Expected behavior |
|---|---|
| Document missing @timestamp | Rejected; producer/DLQ records the failure. |
| Unknown field under strict mapping | Rejected; schema change requires reviewed template version. |
| Rollover during active writes | Logical target remains stable; all acknowledged test events are searchable afterward. |
| Late event within accepted policy | Indexed and queryable by its event timestamp. |
| Stale correction OCC metadata | 409 conflict; correction job re-reads rather than blind overwrite. |
| Snapshot restore drill | Restored data is searchable and stream/template/lifecycle state is explicitly validated. |
8. Security and multi-tenant constraints
Telemetry often contains URLs, user identifiers, headers and stack traces. Redact secrets before indexing, encrypt transport, use least-privilege ingest/read roles, and ensure tenant-scoped users cannot search another tenant's documents or backing indices. Backing-index access can matter for administrative operations and PIT/search permissions, so test the real security plugin/license/deployment rather than assuming data-stream permissions automatically imply every backing-index permission.
9. Release gate
schema/template simulation: PASS
data stream creation: PASS
required timestamp enforcement: PASS
rollover cross-generation search: PASS
exact fixture count/sums: PASS
late-event policy fixture: PASS
unknown-label rejection: PASS
correction OCC conflict test: PASS
tenant authorization tests: PASS
indexing lag p95: MEASURED AGAINST SLO
search p95/p99: MEASURED AGAINST SLO
storage bytes/day: MEASURED
snapshot restore drill: PASS
retention policy review: APPROVED
Do not replace MEASURED with invented numbers. Run
the load and restore tests in an environment representative
enough for the intended decision.
Check your understanding
- Why must template simulation remain in the Chapter 10 deployment?
- What is the shared application-level lifecycle contract?
- Why is one shard/zero replicas not a production recommendation?
- What should a rollover acceptance test prove?
- What makes retention complete?
Review the answers
1. Future backing indices are created from templates, so drift may surface only after rollover.
2. Retention/rollover/restore intent and SLOs, implemented with product-specific lifecycle APIs.
3. It is only the reproducible single-node lab topology; production needs availability, throughput, recovery and storage measurements.
4. The logical stream remains stable and acknowledged events remain searchable across physical generations.
5. Live deletion rules plus snapshot/archive/legal-hold/restore policy and verification.
Production judgment
Choose data streams when append-heavy semantics, time-based query patterns and lifecycle partitioning match the workload. Budget high-cardinality labels, monitor lag and storage growth, and keep correction operations exceptional and auditable. Separate telemetry clusters or workloads when analytical scans compete materially with latency-sensitive search. Validate cost and SLOs under representative ingestion/search concurrency.
Summary and next step
Chapter 10 turns the stream into an operational contract: stable logical naming, bounded backing generations, measured rollover/retention, governed schema, explicit late-data/correction policy, security and restore evidence. Chapter 11 moves upstream into ingest pipelines, processors, enrichment, failure handling and data quality.
Authoritative references
- Elastic data streams — Logical stream, hidden backing indices, write index, @timestamp and rollover behavior.
- Elastic use a data stream — Indexing, searching, rollover, and document updates/deletes through backing indices.
- Elastic time series data streams — Elastic-specific TSDS dimensions, metrics and generated time-series identity.
- Elastic time series index settings — index.mode=time_series, routing path, look-back/look-ahead windows and related settings.
- Elastic time-bound indices and dimension routing — Timestamp acceptance windows and dimension-based shard routing.
- Elastic data stream lifecycle — Built-in lifecycle, rollover, retention, downsampling and storage transitions.
- Elastic data stream retention — Effective retention semantics and lifecycle APIs.
- OpenSearch data streams — Backing indexes, timestamp field, rollover, search and ISM integration.
- OpenSearch rollover API — Manual rollover semantics for data streams and aliases.
- OpenSearch modify data stream API — OpenSearch 3.8 experimental backing-index add/remove operation.
- OpenSearch Index State Management — Policy-driven rollover, retention and index-state automation.
- OpenSearch update document API — Document correction semantics when targeting a concrete backing index.