AWS 304: Amazon MSK and OpenSearch Service
Why this lesson matters
Amazon Managed Streaming for Apache Kafka (Amazon MSK) and Amazon OpenSearch Service are often placed in one diagram, but they solve different problems. MSK is a managed Apache Kafka data plane for ordered partition logs, consumer groups, replay, and Kafka ecosystem compatibility. OpenSearch is a distributed search and analytics engine for documents, logs, metrics, and vectors. Kafka is not a search index. OpenSearch is not a durable event backbone or system of record.
Combining them can be valuable: applications publish immutable events to Kafka, an ingestion component transforms them, and OpenSearch provides operational search. The combination also creates two partitioning systems, two retention models, two scaling planes, two security models, and an offset-versus-index-commit consistency problem. A green cluster on each side does not prove that every accepted event became one searchable document.
This workshop starts with Kafka and OpenSearch internals, compares current managed deployment choices, and then joins them through an explicit ingestion contract. You will design capacity, identities, networks, schemas, failure behavior, replay, recovery, lifecycle, cost, and cleanup for an enterprise event-and-search platform. The required path creates no resources because both services can generate meaningful cost while idle.
What you will be able to do
By the end, you can:
- explain Kafka brokers, topics, partitions, leaders, followers, replicas, in-sync replicas, offsets, producers, and consumer groups;
- compare MSK Provisioned Standard brokers, Provisioned Express brokers, and MSK Serverless;
- size broker, partition, storage, network, and consumer capacity from peak evidence;
- configure producer acknowledgements, idempotence, replication, retention, and failure safety;
- choose IAM, SASL/SCRAM, mutual TLS, and network connectivity patterns;
- explain OpenSearch documents, indexes, mappings, analyzers, primary shards, replicas, cluster-manager nodes, data tiers, and snapshots;
- compare provisioned domains with OpenSearch Serverless search, time-series, and vector collections;
- design index templates, rollover, lifecycle, shard count, refresh, and mappings without a shard or field explosion;
- select MSK Connect, OpenSearch Ingestion, Flink, Lambda, or custom consumers for a justified ingestion path;
- reconcile Kafka offsets, retries, document IDs, bulk writes, and failed records;
- diagnose broker, partition, consumer-lag, cluster-health, indexing, JVM, storage, and query failures; and
- produce an auditable service-selection, recovery, and cost dossier.
Before you start
- Complete AWS303 for streaming semantics, replay, duplicates, backpressure, and Flink.
- Use a non-production evidence path unless budget, quotas, VPC, identities, logs, and cleanup have explicit approval.
- The example Region is
ap-south-1; verify exact broker type, Kafka version, collection type, instance, and integration support. - Do not run Kafka admin commands against an unknown cluster or OpenSearch writes against an unknown domain.
- Do not enable public endpoints merely to simplify a lesson connection.
- Redact bootstrap brokers, VPC endpoints, account IDs, domain endpoints, usernames, certificates, topic contents, and signed requests.
- Treat source events and indexed documents as potentially sensitive, even when the search UI looks harmless.
Create a local evidence directory:
mkdir -p "$HOME/aws304-evidence"
cd "$HOME/aws304-evidence"
export AWS_DEFAULT_REGION="ap-south-1"
aws sts get-caller-identity --query Arn --output text
aws configure list
1. Build the Kafka mental model
producer key and value
|
v
topic partition leader on one broker
|
+--> follower replica on another broker/AZ
+--> follower replica on another broker/AZ
|
v
ordered offsets within that partition
|
+--> consumer group A: one active consumer per assigned partition
+--> consumer group B: independent offset and replay
A topic is split into partitions. Ordering exists within a partition, not globally across the topic. A keyed producer normally hashes the key to a partition so events for the same entity retain partition order. A partition has one leader that handles reads/writes and follower replicas. The in-sync replica set is the replicas sufficiently caught up to participate safely.
An offset identifies a record position in one partition. A consumer group divides partitions among its members. More consumers than partitions leave some consumers idle. Adding partitions can increase parallelism, but does not retroactively repartition old data and can change key-to-partition mapping for future records.
Kafka retention is independent of whether a consumer read a record. Time- or size-based cleanup makes Kafka a retained log, not a queue that deletes after acknowledgement. Log compaction retains the latest value per key according to Kafka semantics; it is not arbitrary deduplication or a substitute for business history.
| Concept | Design question |
|---|---|
| Topic | Which event contract and owner? |
| Partition key | What ordering domain and load distribution? |
| Partition count | What peak parallelism and broker overhead? |
| Replication factor | How many broker copies? |
min.insync.replicas | How many replicas must remain eligible for safe writes? |
Producer acks | What acknowledgement durability is required? |
| Consumer offset | When is progress committed relative to side effects? |
| Retention | How long can every required consumer or replay lag? |
2. Compare current Amazon MSK deployment choices
| Choice | Control model | Strength | Tradeoff |
|---|---|---|---|
| MSK Provisioned with Standard brokers | Choose broker size/count, EBS, Kafka config | Broad Kafka control and mature feature set | Storage, brokers, partition density, and scaling need ownership |
| MSK Provisioned with Express brokers | Choose broker size/count; managed elastic storage and enforced defaults | Higher throughput per broker, faster scale/recovery, less storage management | Supported features/configuration differ; client quotas and sizing still apply |
| MSK Serverless | Service manages broker capacity; pay for throughput/storage dimensions | Lowest broker administration for variable compatible workloads | Less broker-level control, feature/throughput/quota/authentication boundaries |
Express is a broker type within MSK Provisioned, not the same as MSK Serverless. Current Express brokers use managed pay-as-you-go storage, best-practice availability defaults, and client throughput quotas. Standard brokers remain appropriate when a required Kafka configuration or feature is not supported on Express.
MSK Serverless is not automatically cheapest. Compare steady high throughput, multiple consumers, cross-AZ traffic, storage, and predictability. Verify Kafka client and feature compatibility before choosing it for an existing platform.
Write an ADR that rejects at least one mode for a concrete reason. “Managed” and “serverless” are not measurable requirements.
3. Size MSK from every traffic path
Broker ingress is not the entire load:
total broker work includes
producer ingress
+ replication traffic
+ each consumer group's egress
+ catch-up and replay traffic
+ partition movement and broker recovery
+ protocol, compression, and request overhead
Build a workload sheet with average and p99:
- messages and uncompressed/compressed bytes per second;
- batch size, compression codec, acknowledgement mode, and producer concurrency;
- partition-key distribution and hottest key;
- consumer groups and bytes each reads;
- retention, replication factor, and daily storage growth;
- replay and disaster-recovery traffic;
- connection counts and request rates; and
- growth over the capacity-change lead time.
Partition count must satisfy parallelism and throughput without overwhelming brokers with metadata, file handles, elections, replication, and recovery time. Too few partitions cause hot leaders and limited consumers. Too many increase overhead and make broker replacement or reassignment slower.
For Standard brokers, size EBS capacity and throughput. Storage auto scaling can increase capacity but does not shrink it, fix poor retention, or remove disk-throughput bottlenecks. For Express, elastic storage removes volume provisioning but not retention cost or partition management.
4. Design Kafka durability and client safety
A strong baseline for critical events commonly includes replication across three Availability Zones, replication factor 3, min.insync.replicas=2, producer acks=all, idempotent production, bounded retries, and unclean leader election disabled. Validate exact broker-mode defaults and compatibility rather than copying settings blindly.
producer sends event with stable event_id
|
leader appends and required in-sync replicas acknowledge
|
producer receives success
|
consumer reads and performs idempotent/transactional side effect
|
consumer commits offset only after safe outcome
Producer idempotence prevents certain duplicates caused by producer retries within its supported session and configuration. It does not deduplicate a business event published again by another process. Kafka transactions can coordinate supported Kafka writes and offsets, but do not make an arbitrary OpenSearch, email, or payment side effect atomic.
Consumers must decide between at-most-once risk, at-least-once replay, and a truly transactional integration. For an OpenSearch sink, use a deterministic document ID derived from event_id or an entity/version key where overwriting is correct. Preserve failures separately. Commit the Kafka offset only according to the connector's proven delivery semantics.
5. Design topics, schemas, and lifecycle
Each topic needs:
- owner, classification, producer and consumer contracts;
- key/value schema and compatibility rule;
- partition count and change procedure;
- retention/compaction and legal hold;
- replication and minimum ISR;
- maximum record size and compression;
- access-control groups and network path;
- lag/freshness SLOs;
- replay and dead-letter strategy; and
- deprecation/deletion approval.
Use a schema registry when multiple producers and consumers need governed evolution. Compatibility mode is not sufficient alone; semantic changes such as currency units or event meaning can remain syntactically compatible and still break consumers.
Avoid putting unbounded payloads in Kafka. Store large immutable objects in S3 and publish a versioned reference with checksum, ownership, expiry, and authorization, if that pattern meets reliability needs.
6. Build MSK identity and network paths
MSK supports combinations of TLS encryption and authentication options that vary by cluster type and configuration, including IAM access control, SASL/SCRAM backed by Secrets Manager, and mutual TLS. Choose from client ecosystem and identity lifecycle, not convenience.
IAM authentication integrates with AWS roles but requires supported client libraries and correctly signed connections. SCRAM needs secret rotation and broker association. Mutual TLS needs private CA, certificate issuance, revocation, renewal, and client-keystore operations. Encrypt client traffic and data at rest with an owned KMS design.
Draw every path:
producer/consumer subnets
-> DNS and MSK bootstrap brokers
-> broker security groups and listener port
-> authentication dependency: IAM/STS, secret, or certificate
-> schema registry and monitoring
-> connector worker and destination
Multi-VPC private connectivity, PrivateLink patterns, peering, transit routing, and public access have different support and cost. Include route tables, security groups, NACLs, DNS, cross-AZ data, NAT or endpoints, and client source addresses. A cluster state of ACTIVE does not prove any client can resolve, authenticate, or reach brokers.
7. Operate Kafka as a distributed system
Monitor:
- offline partitions and active-controller health;
- under-replicated partitions and ISR changes;
- broker CPU, memory, network, storage, and request latency;
- produce/fetch throttling and errors;
- partition count, leader distribution, and skew;
- consumer lag by group and partition;
- connection/authentication failures;
- replication and reassignment progress; and
- broker logs and client errors.
Consumer lag in records is not a time SLO by itself. Convert it using current input and sustainable processing rates. A group with 10 million tiny events differs from one with 10 million large events.
Maintenance and version upgrades require compatibility testing for clients, protocol, message format, connectors, and configurations. Test broker reboot at peak load. Clients must refresh metadata and fail over. Capacity should tolerate one broker unavailable without violating latency or write safety.
Use MSK Replicator or another proven replication approach only after defining topic selection, naming, offset synchronization, loop prevention, failover authority, RPO/RTO, and failback. Multi-Region replication is not automatic active-active conflict resolution.
8. Build the OpenSearch mental model
JSON document
|
mapping and analyzer
|
index primary shard ----> replica shard on another node/AZ
|
distributed inverted index / vector structures
|
query coordinator fans out, merges results, returns response
An index contains documents. Mappings define field types. Text analyzers tokenize full-text fields; keyword fields preserve exact values for filters, aggregations, and sorting. A primary shard owns part of an index; replicas improve availability and read capacity. Shard count is chosen at index creation and has lasting operational consequences.
OpenSearch is near-real-time: a successful index response and a refresh boundary are different events. Search results can lag accepted writes. A refresh every second, high replica count, and tiny bulk requests can destroy indexing throughput.
Do not use OpenSearch as the only copy of orders, payments, identities, or compliance records. Rebuildability requires an authoritative source, schema/templates, ingestion position, and replay runbook.
9. Compare provisioned domains and OpenSearch Serverless
| Choice | Unit | Best fit | Main ownership |
|---|---|---|---|
| Provisioned domain | Domain with cluster-manager, data, warm/cold/coordinator choices | Predictable sustained workloads, deep tuning, controlled topology | nodes, shards, tiers, storage, versions, updates |
| Serverless search collection | OCUs and managed storage | full-text/application search with variable demand | collection/security policies, capacity bounds, mappings |
| Serverless time-series collection | OCUs, hot/warm behavior, lifecycle | logs and time-oriented analytics | retention/lifecycle, ingest/query capacity |
| Serverless vector collection | OCUs and vector indexes | semantic/vector search | dimensions, method, memory, recall/latency/cost |
OpenSearch Serverless decouples managed compute and storage and measures capacity in OpenSearch Compute Units. Current collection groups can set minimum and maximum indexing and search OCUs and share capacity within defined boundaries. Scaling to zero can reduce idle compute but adds cold-start implications. A maximum protects budget and can also throttle workload.
Collection type is not cosmetic and cannot be mixed within one collection group. Time-series lifecycle behavior differs from search and vector collections. Validate current features and API compatibility before migrating a domain workload.
10. Design provisioned OpenSearch domains
For production domains, evaluate Multi-AZ with Standby. It enforces three-AZ and topology practices, maintains standby capacity, and avoids recovery redistribution for an Availability Zone failure. Standard Multi-AZ without standby can redistribute shards under failure, increasing pressure exactly when capacity is reduced.
Separate roles:
- dedicated cluster-manager nodes maintain cluster state and elections;
- data nodes index, store, and search;
- dedicated coordinator nodes can offload request coordination for justified workloads;
- UltraWarm and cold storage can reduce older-data cost with different latency; and
- ingest pipelines or OpenSearch Ingestion transform before indexing.
Size heap, CPU, storage, IOPS/throughput, network, shard count/size, replicas, indexing, query concurrency, aggregations, and recovery headroom. Keep free storage above operational watermarks. A domain at 70 percent normal utilization may have no safe room for node loss, shard movement, or a traffic spike.
Blue/green service updates can temporarily need extra capacity and can alter endpoints or behavior according to the change. Use off-peak windows, maintenance evidence, client retries, compatibility tests, and rollback planning.
11. Govern mappings, shards, and index lifecycle
Dynamic mappings can create thousands of fields from uncontrolled JSON keys. Mapping explosion consumes cluster state and heap. Explicit templates should define:
- timestamp and document ID;
textversuskeyword;- numeric/date/IP/geo/vector types;
- dynamic-field policy;
- sensitive fields excluded or transformed;
- analyzers and normalizers;
- primary/replica count;
- refresh interval; and
- rollover and lifecycle policy.
Time-based data should use rollover based on size, age, or document count rather than hard-coded daily indexes when volumes vary. Index State Management can transition or delete indexes for provisioned domains. Serverless time-series data lifecycle uses its own supported policy model.
Shard sizing is workload-specific. Too many small shards consume heap and coordination. Huge shards slow recovery and movement. Validate with ingestion/query/recovery tests; do not use one universal shard-size number as a certification fact.
Snapshots are recovery artifacts, not high availability. OpenSearch Service takes automated snapshots according to the service model, and manual repository approaches have separate requirements. Test restore into an isolated target, validate mappings/security/lifecycle, and reconcile writes after the snapshot point.
12. Join MSK to OpenSearch without losing evidence
Possible paths include MSK Connect with a compatible sink connector, OpenSearch Ingestion where the supported source and sink match, Managed Flink, Lambda for bounded patterns, or a custom consumer. Compare:
| Path | Strength | Risk |
|---|---|---|
| MSK Connect | Managed Kafka Connect workers and ecosystem | connector quality, offset/commit semantics, scaling |
| OpenSearch Ingestion | Managed pipeline and OpenSearch-oriented processing | source/processor feature fit and OCU cost |
| Managed Flink | stateful transforms, event time, complex routing | application/state/connector operations |
| Lambda | simple low-volume event handling | batch retry, duration, concurrency, backpressure |
| Custom consumer | maximum behavior control | full ownership of leases, retries, deployment, recovery |
The ingestion contract must include:
Kafka topic/partition/offset/event_id
|
deserialize and validate schema
|
transform to versioned OpenSearch document
|
bulk index with deterministic document ID
|
classify every item response
| success | retryable | permanent
|
commit offset only under proven connector semantics
|
reconcile source events, successful docs, retries, and quarantine
OpenSearch bulk HTTP success does not mean every item succeeded. Parse item-level results. Retry only transient failures. Quarantine permanent mapping/validation failures with original topic, partition, offset, event ID, error, and schema version.
Replay into the same deterministic IDs can rebuild state when overwrite semantics are correct. Append-only search history may require a unique event ID. Delete/tombstone semantics and out-of-order entity versions need explicit conflict control.
13. Secure OpenSearch access
Provisioned domain authorization can involve network reachability, domain access policy, IAM request signing, and fine-grained access control. OpenSearch Serverless separately uses encryption, network, and data-access policies. Passing one layer does not bypass another.
Use VPC endpoints for private workloads unless a justified public pattern is secured. OpenSearch Dashboards access needs identity federation or an approved authentication flow, least-privilege roles, and audit logging. Never share a master-user credential.
Field- and document-level security can restrict query results but does not transform source data or erase exported results. KMS policies, snapshots, logs, ingestion buffers, failed records, and downstream dashboards require their own controls.
Protect clusters from abusive queries with role controls, bounded result windows, aggregation limits, timeouts, circuit breakers, and tenant separation. Search access is resource consumption.
14. Diagnose from both sides of the pipeline
| Symptom | Strong evidence | Likely direction |
|---|---|---|
MSK ACTIVE, clients cannot connect | DNS, route, SG, listener, auth logs | network or authentication path |
| Producer latency and under-replicated partitions | ISR, broker/network/storage metrics | broker overload or replica catch-up |
| One broker hot | partition leaders and key distribution | partition/key/leader skew |
| Consumer group constantly rebalances | client logs, session/poll timing | slow processing or unstable membership |
| Disk rises despite retention | segment/retention config and traffic | policy mismatch, lag, storage throughput |
| OpenSearch cluster red | unassigned primary shards and allocation reason | unavailable primary/capacity/allocation |
| Cluster yellow | unassigned replicas | topology or allocation capacity |
| Bulk request returns errors | item-level response | mapping, rejection, version conflict |
| Write rejections and high JVM | thread pools, heap, GC, shards | overload, mappings, shard design |
| Search latency high | slow logs, query profile, cache, shard fan-out | expensive query or too many shards |
| Kafka lag low but documents missing | connector quarantine and reconciliation | permanent item failures or bad commit |
| Replay duplicates documents | document-ID and version strategy | non-idempotent sink design |
Do not respond by adding brokers and OpenSearch nodes simultaneously. Locate whether source, connector, destination, or query load is limiting, then change one controlled variable.
15. P18 architecture exercise
Design a platform for:
- 25,000 events/s normal and 100,000 events/s campaign peak;
- 1.5 KiB compressed average message;
- six producer services and four independent consumer groups;
- seven-day Kafka replay and one-year S3 archive;
- searchable operational events within 30 seconds;
- 30-day hot log search and 180-day lower-cost retention;
- deterministic rebuild of indexes;
- regulated tenant isolation;
- one-Region RPO 0 for an AZ failure and documented Regional DR; and
- no public broker or search endpoint.
Produce MSK mode/broker/partition/storage/authentication design, topic/schema/retention catalog, producer/consumer safety, OpenSearch domain-versus-Serverless ADR, index template and lifecycle, ingestion commit/retry/quarantine model, reconciliation, scaling, Multi-AZ behavior, snapshots/rebuild, Regional recovery, monitoring, quotas, cost, and cleanup.
Include failure games for one broker unavailable, hot partition, consumer lag beyond retention forecast, connector poison event, destination write rejection, OpenSearch node/AZ failure, red index, KMS denial, expired certificate/secret, schema incompatibility, snapshot restore, and full index rebuild from Kafka/S3.
16. Cost, quotas, inventory, and cleanup
MSK cost can include broker or Serverless capacity, storage, storage throughput, data transfer, multi-VPC connectivity, public IPv4 where applicable, MSK Connect workers, Replicator, logs, monitoring, Secrets Manager, Private CA, KMS, and cross-Region replication. OpenSearch cost can include domain instances or Serverless OCUs, EBS/IOPS, UltraWarm/cold storage, snapshots, OpenSearch Ingestion OCUs, transfer, audit logs, KMS, NAT, and dashboards.
Read-only inventory:
aws kafka list-clusters-v2 --output json | tee msk-clusters.json
aws kafka list-serverless-clusters --output json | tee msk-serverless.json
aws kafka list-replicators --output json | tee msk-replicators.json
aws opensearch list-domain-names --output json | tee opensearch-domains.json
aws opensearchserverless list-collections --output json | tee opensearch-collections.json
Cleanup for an approved lab proves cluster deletion, connector/replicator removal, topic/export retention decisions, broker ENIs and security groups handled, domain/collection deletion, ingestion pipelines stopped/deleted, snapshots retained or removed by policy, and cleanup of S3, logs, KMS grants, secrets/certificates, IAM roles, alarms, dashboards, VPC endpoints, and DNS. Access denied is not deletion proof. Review delayed billing.
17. Practical submission
Submit a p18-msk-opensearch/ package containing:
- separate Kafka and search requirement statements;
- MSK Standard/Express/Serverless ADR;
- peak ingress, replication, egress, storage, partition, and broker model;
- topic/key/schema/retention/compaction ownership catalog;
- producer durability and consumer offset/side-effect state machines;
- authentication, KMS, VPC, DNS, and client path;
- maintenance, version, rebalance, and Regional recovery runbooks;
- OpenSearch domain-versus-collection ADR;
- node/OCU, shard/replica, storage/tier, and failure-headroom model;
- explicit mappings, templates, rollover, lifecycle, and snapshot plan;
- ingestion component ADR and item-level error contract;
- offset-to-document reconciliation and deterministic rebuild;
- twelve failure-injection evidence packs;
- monitoring tied to lag, freshness, durability, and search SLOs;
- complete cost/quota sensitivity; and
- cleanup and delayed billing evidence.
Knowledge check
- Where does Kafka guarantee order? Within one partition.
- Why can adding partitions affect keyed order? Future key hashing can map to a different partition.
- Is Express the same as MSK Serverless? No; Express is an MSK Provisioned broker type.
- Does producer idempotence deduplicate every business duplicate? No.
- What does a consumer offset prove? A position chosen by that group, not atomic external completion.
- Why is OpenSearch not a source of truth? Indexes can be transformed, deleted, reindexed, or temporarily unavailable.
- What is the difference between primary and replica shards? Primaries own index partitions; replicas provide redundant copies and reads.
- Is OpenSearch Serverless one generic collection type? No; search, time-series, and vector choices differ.
- Why parse every bulk item response? The HTTP request can succeed while individual documents fail.
- What makes index replay safe? Deterministic IDs/version rules, quarantine, offset evidence, and reconciliation.
Lesson acceptance
Pass only when all are true:
- Kafka partitions, replicas, ISR, acknowledgements, offsets, groups, retention, and compaction are accurately separated;
- Standard, Express, and Serverless MSK are compared from measured requirements;
- capacity includes producer, replication, every consumer, replay, storage, and recovery;
- producer and consumer semantics protect ordering and duplicate-safe side effects;
- topic/schema/retention/access ownership is explicit;
- authentication, KMS, private connectivity, and client failover are testable;
- OpenSearch mappings, analyzers, shards, replicas, nodes/OCUs, tiers, lifecycle, and snapshots are designed;
- provisioned domains and all relevant Serverless collection types are evaluated;
- OpenSearch remains rebuildable from an authoritative source;
- ingestion handles each bulk item, deterministic IDs, retry classes, quarantine, and offset commits;
- Kafka-to-index reconciliation detects missing and duplicate documents;
- twelve failure cases cover source, ingestion, destination, security, and recovery;
- cost includes both idle and variable capacity plus all connected services; and
- cleanup handles compute, storage, connectors, endpoints, identities, logs, snapshots, and delayed charges.
Official sources
- What is Amazon MSK?
- MSK broker types
- MSK Express brokers
- MSK Express best practices
- MSK authentication and authorization
- Troubleshoot Amazon MSK
- What is Amazon OpenSearch Service?
- Create and manage OpenSearch domains
- OpenSearch Multi-AZ with Standby
- Amazon OpenSearch Serverless
- OpenSearch Serverless capacity
- OpenSearch cost optimization