AWS courseLesson 11 of 14
AWS course · Lesson 11 of 14
AWS Data Services Landscape: Networking, Databases, Messaging, CDC and Cost
The AWS services around a data platform: VPC networking, RDS and Aurora, DynamoDB, SQS and SNS, DMS for CDC, Amazon MSK, Cost Explorer and Budgets, Well-Architected.
On this page
A data platform on AWS is more than S3, Glue and a warehouse. Pipelines run inside networks, read from operational databases, pass messages through queues, replicate changes with DMS, sometimes stream through Kafka, and all of it costs money that someone will ask you to explain. This lesson gives you working knowledge of these surrounding services: enough to design with them, debug the common failures and answer interview questions confidently.
Sample code
The runnable examples are a subnet calculation, a model of SQS redelivery and dead-lettering, and a PostgreSQL merge that applies change records the way a CDC pipeline does. Commands that need an AWS account are shown but not executed.
VPC and networking basics
What it is. A VPC (virtual private cloud) is your private network in an AWS Region, with an IP range (CIDR block) you choose. Data services that run “in your VPC” (RDS, Redshift, EMR, MSK, Glue jobs with connections, Lambda functions attached to a VPC) get network interfaces in its subnets.
How it works.
| Building block | What it does | Data engineering relevance |
|---|---|---|
| Subnet | A slice of the VPC CIDR in one Availability Zone | Spread databases and clusters over at least two AZs |
| Route table | Where traffic from a subnet goes | Decides if a subnet is public or private |
| Internet gateway | Route to the internet for public subnets | Rarely needed by data services directly |
| NAT gateway | Lets private subnets reach the internet outbound | Billed per hour and per GB processed; large S3 traffic through NAT is a classic cost leak |
| Security group | Stateful allow-list on network interfaces | Allow Glue or EMR to reach the database port; Glue needs a self-referencing rule |
| Network ACL | Stateless allow and deny rules per subnet | Coarse guardrails |
| Gateway VPC endpoint | Private route to S3 or DynamoDB, no charge for the endpoint | Keeps lake traffic off NAT and the internet |
| Interface VPC endpoint (PrivateLink) | Private IP for other AWS services (Glue, STS, Secrets Manager, KMS) | Needed when private subnets have no NAT |
Every subnet loses five addresses to AWS, so size subnets for EMR clusters, Glue workers and Lambda functions, which all take IPs:
import ipaddress
vpc = ipaddress.ip_network("10.20.0.0/16")
subnets = list(vpc.subnets(new_prefix=20))[:6]
names = ["public-a", "public-b", "private-a", "private-b", "data-a", "data-b"]
for name, net in zip(names, subnets):
usable = net.num_addresses - 5 # AWS reserves 5 addresses in every subnet
print(f"{name:<10} {str(net):<16} {usable:>5} usable IPs")
public-a 10.20.0.0/20 4091 usable IPs
public-b 10.20.16.0/20 4091 usable IPs
private-a 10.20.32.0/20 4091 usable IPs
private-b 10.20.48.0/20 4091 usable IPs
data-a 10.20.64.0/20 4091 usable IPs
data-b 10.20.80.0/20 4091 usable IPs
aws ec2 create-vpc-endpoint --vpc-id vpc-0abc1234def567890 \
--service-name com.amazonaws.eu-west-1.s3 --vpc-endpoint-type Gateway \
--route-table-ids rtb-0private1 rtb-0private2
Pitfalls.
- A Glue job or Lambda function in a private subnet that times out reaching S3 or STS: there is no NAT gateway or VPC endpoint. Add the S3 gateway endpoint and interface endpoints for the other services.
- Subnets too small for a large EMR cluster or a Glue job with many workers, which fail with “not enough IP addresses”.
- Security groups that allow
0.0.0.0/0to a database port “to make it work”.
In interviews. Draw a VPC with public and private subnets across two AZs, databases and clusters in private subnets, an S3 gateway endpoint, and explain why NAT gateway data processing charges matter for data platforms.
RDS and Aurora basics
What it is. Amazon RDS runs managed relational databases (PostgreSQL, MySQL, MariaDB, Oracle, SQL Server, Db2). Amazon Aurora is AWS’s MySQL- and PostgreSQL-compatible engine with a distributed storage layer. For data engineers they are usually sources: the operational systems whose data you extract.
How it works.
- Multi-AZ keeps a synchronous standby in another AZ for failover; it is for availability, not for reading.
- Read replicas are asynchronous copies that serve reads. Point extract queries and CDC-unfriendly bulk reads at a replica so analytics never slows the production primary.
- Aurora stores six copies of data across three AZs, supports up to 15 low-lag replicas, and has Aurora Serverless v2 for variable load.
- Snapshot export to S3 writes a database snapshot as Parquet without touching the live database.
- Zero-ETL integrations replicate Aurora (and some RDS engines) into Redshift continuously without pipelines you build.
- For CDC, the source must expose its change log: binary logging in row format (MySQL), or logical replication enabled (PostgreSQL,
rds.logical_replication), plus retained backups or logs long enough for the replication tool to keep up.
Pitfalls.
- Full-table
SELECT *extracts from the primary at business hours. - Incremental extracts on an
updated_atcolumn that misses deletes and rows updated without the timestamp changing. Use CDC for correctness. - Replication slots on PostgreSQL that are not consumed, which make the database retain WAL until the disk fills.
In interviews. Know Multi-AZ versus read replicas, the options for getting data out (replica queries, snapshot export, DMS CDC, zero-ETL) and the source settings CDC needs.
DynamoDB for data engineers
What it is. DynamoDB is a serverless key-value and document database with single-digit millisecond reads and writes at any scale. Data engineers meet it as a source to extract, as a fast lookup or state store for pipelines (processed-file registers, idempotency keys, watermarks), and as a serving layer for precomputed results.
How it works.
- Items are addressed by a partition key (hashed to spread data) and an optional sort key (ordering within a partition). Items can be up to 400 KB.
- Design starts from access patterns: there are no joins, and queries need a key. Global secondary indexes provide other access paths.
- Capacity is on-demand (pay per request) or provisioned (with auto scaling).
- DynamoDB Streams record item-level changes for 24 hours; Lambda can process them. Kinesis Data Streams for DynamoDB is the alternative for longer retention and more consumers.
- Export to S3 (full or incremental, which requires point-in-time recovery) writes table data to S3 without consuming read capacity; this is the right way to bring DynamoDB data into the lake.
- TTL deletes expired items automatically, useful for idempotency keys.
import boto3
from botocore.exceptions import ClientError
# Table key: partition key object_key, sort key etag
table = boto3.resource("dynamodb").Table("processed_files")
def claim(object_key, etag):
"""Return True the first time a file version is seen; False for duplicates."""
try:
table.put_item(
Item={"object_key": object_key, "etag": etag},
ConditionExpression="attribute_not_exists(object_key)", # no item with this key yet
)
return True
except ClientError as e:
if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
return False
raise
Pitfalls.
- Scanning a large table for analytics, which consumes capacity and is slow. Use export to S3.
- Hot partitions from a low-cardinality partition key (for example
date). - Treating DynamoDB like a relational database and needing joins later.
In interviews. Explain partition and sort keys, access-pattern-first design, Streams versus export to S3, and the conditional-write idempotency pattern above.
SQS and SNS
What it is. Amazon SQS is a message queue: producers send messages, consumers poll, process and delete them. Amazon SNS is publish/subscribe: a message published to a topic is pushed to every subscriber (SQS queues, Lambda functions, HTTP endpoints, email). EventBridge is an event bus with content-based routing and many AWS service events.
How it works.
| SQS Standard | SQS FIFO | SNS | |
|---|---|---|---|
| Delivery | At least once | Exactly-once processing within the deduplication interval | Push to all subscribers |
| Ordering | Best effort | Per message group | FIFO topics only |
| Throughput | Very high | Lower (higher with high-throughput mode) | High |
| Use | Buffer work, decouple producers and consumers | Ordered processing per key | Fan-out one event to many consumers |
Key SQS settings:
- Visibility timeout (default 30 seconds, up to 12 hours): after a consumer receives a message, it is hidden for this long. If the consumer does not delete it in time, it reappears and is processed again.
- Retention: from 1 minute to 14 days (default 4 days).
- Message size: up to 1 MiB (raised from 256 KiB in August 2025); larger payloads go to S3 with a pointer, for example with the extended client libraries. SNS also supports 1 MiB payloads since September 2026.
- Long polling (
WaitTimeSecondsup to 20) reduces empty receives and cost. - Dead-letter queue: a redrive policy moves a message to a DLQ after
maxReceiveCountreceives, so one poison message cannot loop forever. You can redrive messages from the DLQ back to the source queue once fixed.
This model shows a visibility timeout of 30 seconds and maxReceiveCount of 3: file-2 fails once and succeeds on redelivery, file-3 always fails and lands in the DLQ:
def sqs_consume(messages, fails_until, visibility_timeout=30, max_receive_count=3, poll_every=10, horizon=200):
"""Model an SQS queue with a redrive policy. fails_until[msg] = attempts that fail before success."""
visible_at = {m: 0 for m in messages}
receives = {m: 0 for m in messages}
done, dlq, log = [], [], []
for t in range(0, horizon, poll_every):
for m in sorted(visible_at):
if visible_at[m] > t:
continue
if receives[m] >= max_receive_count: # redrive happens on the next receive attempt
dlq.append(m); del visible_at[m]; log.append(f"t={t:>3} {m} -> DLQ")
continue
receives[m] += 1
if receives[m] > fails_until.get(m, 0):
done.append(m); del visible_at[m]; log.append(f"t={t:>3} {m} processed (receive {receives[m]})")
else:
visible_at[m] = t + visibility_timeout # not deleted: reappears after the timeout
return done, dlq, log
done, dlq, log = sqs_consume(["file-1", "file-2", "file-3"], {"file-2": 1, "file-3": 99})
print("\n".join(log))
print("processed:", done, "| dead-letter queue:", dlq)
t= 0 file-1 processed (receive 1)
t= 30 file-2 processed (receive 2)
t= 90 file-3 -> DLQ
processed: ['file-1', 'file-2'] | dead-letter queue: ['file-3']
{
"RedrivePolicy": "{\"deadLetterTargetArn\":\"arn:aws:sqs:eu-west-1:111122223333:landing-dlq\",\"maxReceiveCount\":\"3\"}",
"VisibilityTimeout": "180",
"MessageRetentionPeriod": "1209600",
"ReceiveMessageWaitTimeSeconds": "20"
}
aws sqs set-queue-attributes \
--queue-url https://sqs.eu-west-1.amazonaws.com/111122223333/landing \
--attributes file://queue-attributes.json
The common fan-out pattern is SNS to several SQS queues: one event (a file landed) reaches the lake loader, the search indexer and the audit logger, each with its own queue, retries and DLQ.
Pitfalls.
- Visibility timeout shorter than processing time, so messages are processed twice in parallel. Set it above the consumer’s maximum time (for Lambda, AWS recommends at least six times the function timeout).
- No DLQ, or a DLQ without an alarm.
- Assuming Standard queues keep order or deliver once.
In interviews. Explain visibility timeout, at-least-once delivery, DLQs with maxReceiveCount, Standard versus FIFO, and SNS-to-SQS fan-out. Use the AWS decision guide’s framing: SQS to buffer, SNS to broadcast, EventBridge to route events by content.
AWS DMS (migration and CDC)
What it is. AWS Database Migration Service copies data between databases and into S3, Redshift, Kinesis, MSK and other targets. Its main data engineering use is change data capture (CDC): a full load of existing rows, then continuous replication of inserts, updates and deletes read from the source’s transaction log.
How it works.
- Components: a replication instance (or DMS Serverless, which scales capacity automatically), a source endpoint, a target endpoint, and a task.
- Task types:
full-load,cdc, orfull-load-and-cdc. - Table mappings (JSON) choose schemas and tables and can rename or transform them.
- With an S3 target, CDC files contain an operation column (
I,UorD) and can include a commit timestamp column; full-load files do not include the operation column unless you configure it. Write Parquet for the lake. - Data validation compares source and target rows to detect drift.
{
"rules": [
{
"rule-type": "selection",
"rule-id": "1",
"rule-name": "include-orders",
"object-locator": { "schema-name": "public", "table-name": "orders" },
"rule-action": "include"
}
]
}
aws dms create-replication-task --replication-task-identifier orders-to-lake \
--source-endpoint-arn arn:aws:dms:eu-west-1:111122223333:endpoint:SRCPG \
--target-endpoint-arn arn:aws:dms:eu-west-1:111122223333:endpoint:S3LAKE \
--replication-instance-arn arn:aws:dms:eu-west-1:111122223333:rep:REPL1 \
--migration-type full-load-and-cdc \
--table-mappings file://table-mappings.json
Applying the changes. DMS lands change records; a downstream job applies them to the current-state table. Keep only the latest change per key (by commit order), then merge: delete on D, update on U, insert on I. This PostgreSQL 16 example is a stand-in for the same MERGE you would run in Athena on an Iceberg table, in Redshift, or in Spark:
CREATE TABLE orders (order_id int PRIMARY KEY, status text, amount numeric(10,2), updated_at timestamp);
INSERT INTO orders VALUES (1, 'placed', 20.00, '2026-10-04 09:00'), (2, 'placed', 35.50, '2026-10-04 10:00');
CREATE TABLE orders_cdc (op char(1), order_id int, status text, amount numeric(10,2), updated_at timestamp, seq bigint);
INSERT INTO orders_cdc VALUES
('U', 1, 'shipped', 20.00, '2026-10-05 08:00', 101),
('I', 3, 'placed', 12.00, '2026-10-05 08:05', 102),
('U', 3, 'paid', 12.00, '2026-10-05 08:07', 103),
('D', 2, NULL, NULL, '2026-10-05 08:10', 104);
MERGE INTO orders AS t
USING (
SELECT DISTINCT ON (order_id) *
FROM orders_cdc
ORDER BY order_id, seq DESC
) AS c
ON t.order_id = c.order_id
WHEN MATCHED AND c.op = 'D' THEN DELETE
WHEN MATCHED THEN UPDATE SET status = c.status, amount = c.amount, updated_at = c.updated_at
WHEN NOT MATCHED AND c.op <> 'D' THEN
INSERT (order_id, status, amount, updated_at) VALUES (c.order_id, c.status, c.amount, c.updated_at);
SELECT * FROM orders ORDER BY order_id;
order_id | status | amount | updated_at
----------+---------+--------+---------------------
1 | shipped | 20.00 | 2026-10-05 08:00:00
3 | paid | 12.00 | 2026-10-05 08:07:00
Order 3 was inserted and updated in the same batch, so only its latest state is applied; order 2’s delete removes it.
Pitfalls.
- Source not prepared (binary logging off, logical replication disabled, logs purged too quickly), so CDC cannot start or falls behind permanently.
- Applying changes without de-duplicating by key and commit order, so an older update overwrites a newer one.
- Large objects (LOB columns) slowing tasks; configure LOB mode deliberately.
- Schema changes on the source: check which DDL the task replicates and how the target handles new columns.
In interviews. Describe full load plus CDC, the operation column, how you apply changes idempotently (latest per key, then merge) and the source prerequisites. Mention monitoring CDC latency.
MSK (managed Kafka)
What it is. Amazon Managed Streaming for Apache Kafka runs Apache Kafka for you: brokers, storage, patching and monitoring. Your producers and consumers use the normal Kafka APIs.
How it works.
- MSK provisioned clusters: you choose broker size and count. Standard brokers give full control of Kafka configuration and storage; Express brokers manage storage for you and scale faster.
- MSK Serverless: no brokers to size; capacity scales automatically, with some limits on configuration and quotas.
- MSK Connect runs Kafka Connect connectors (for example Debezium for CDC, or an S3 sink).
- Authentication with IAM, SASL/SCRAM or mutual TLS; encryption in transit and at rest; tiered storage for long retention on provisioned clusters.
- Firehose can read from MSK and deliver to S3, and Lambda can consume MSK topics with an event source mapping.
| Choose MSK when | Choose Kinesis Data Streams when |
|---|---|
| You already use Kafka clients, Kafka Connect, Kafka Streams or a schema registry | You want the least operational work and native Lambda and Firehose integration |
| You need Kafka features such as transactions and compacted topics | Your throughput fits shard-based scaling or on-demand mode |
| Portability across clouds matters | The team is AWS-only |
Pitfalls. MSK still needs capacity planning (partitions, broker storage and throughput) on provisioned clusters, and partition counts follow the same rules as self-managed Kafka.
In interviews. Map MSK options (provisioned Standard or Express, Serverless, Connect) and compare with Kinesis on operations and ecosystem, linking to the Kafka lessons for partitions and consumer groups.
Cost Explorer and budgets
What it is. Cost Explorer shows and forecasts spend, grouped by service, account, Region, usage type or tag. AWS Budgets alerts when actual or forecast spend or usage crosses a threshold, and can take actions. Cost Anomaly Detection uses machine learning to flag unusual spend. For detailed analysis, Data Exports (the Cost and Usage Report) writes line-item billing data to S3, where you can query it with Athena.
How it works.
- Tag every resource with owner, pipeline and environment, and activate those tags as cost allocation tags in the Billing console (the modern-services lesson covers tag strategy).
- Group Cost Explorer by service, then by tag, to find which pipeline drives spend.
- Create a budget per team or pipeline with alerts at, say, 80% actual and 100% forecast, sent to email or SNS.
- Turn on Cost Anomaly Detection for services such as Glue, Athena and EMR.
aws budgets create-budget --account-id 111122223333 \
--budget '{"BudgetName": "data-platform-monthly", "BudgetLimit": {"Amount": "5000", "Unit": "USD"}, "TimeUnit": "MONTHLY", "BudgetType": "COST", "CostFilters": {"TagKeyValue": ["user:team$data-platform"]}}' \
--notifications-with-subscribers '[{"Notification": {"NotificationType": "FORECASTED", "ComparisonOperator": "GREATER_THAN", "Threshold": 100, "ThresholdType": "PERCENTAGE"}, "Subscribers": [{"SubscriptionType": "EMAIL", "Address": "[email protected]"}]}]'
The budget amount above is an example threshold you would set, not a price. The usual data-platform cost drivers to look for: Athena data scanned, Glue DPU-hours (especially crawlers), idle EMR clusters, NAT gateway data processing, S3 request charges from small files, Redshift clusters running around the clock, and CloudWatch Logs without retention.
Pitfalls. Tags that were never activated (Cost Explorer cannot group by them), and budgets that alert a mailbox nobody reads.
In interviews. Show the loop: tag, activate, explore by tag, budget with forecast alerts, anomaly detection, then fix the top driver.
Well-Architected for data
What it is. The AWS Well-Architected Framework is a set of best practices organised in six pillars. The Data Analytics Lens applies them to data platforms.
| Pillar | For a data platform, ask |
|---|---|
| Operational excellence | Are pipelines in code, tested, deployed by CI/CD, monitored, with runbooks? |
| Security | Least-privilege roles, encryption with KMS, Lake Formation for data access, secrets in Secrets Manager, audit with CloudTrail? |
| Reliability | Idempotent, retryable steps, DLQs, multi-AZ sources, backfill and replay procedures? |
| Performance efficiency | Right engine per workload, columnar formats, partitioning, right-sized compute? |
| Cost optimisation | Serverless where idle time is high, Spot, lifecycle rules, tagging and budgets? |
| Sustainability | Less data scanned and stored, efficient instance types such as Graviton, deleting unused data? |
In interviews. When asked to review or design a platform, walk the six pillars briefly. It shows structure, and each pillar maps to concrete AWS features from this course.
Practice questions
A Glue job in a private subnet hangs when it tries to read from S3. What is wrong?
The subnet has no route to S3: no NAT gateway and no S3 gateway VPC endpoint. Add a gateway endpoint for S3 to the subnet’s route table (and interface endpoints for other services the job calls, such as Glue, STS or Secrets Manager, if there is no NAT). Also check the security group has the self-referencing rule Glue needs.
How would you extract a large DynamoDB table into the lake every day?
Enable point-in-time recovery and use DynamoDB export to S3, full the first time and incremental after that, which does not consume table read capacity. Then convert or merge the export into Parquet or Iceberg with Glue or EMR. For near-real-time changes, use DynamoDB Streams or Kinesis Data Streams for DynamoDB.
Messages in an SQS queue are being processed twice at the same time. Why?
The visibility timeout is shorter than the processing time, so the message reappears and another consumer picks it up while the first is still working. Increase the visibility timeout above the maximum processing time (for Lambda, at least six times the function timeout), extend it during long processing, and make consumers idempotent because Standard queues are at least once anyway.
Design CDC from PostgreSQL on RDS to an Iceberg table in S3.
Enable logical replication on the RDS instance and a sufficient log retention. Run a DMS full-load-and-cdc task to an S3 target in Parquet with the operation and commit timestamp columns. On a schedule, a Glue or EMR job (or Athena MERGE) reads new change files, keeps the latest change per primary key by commit order, and merges into the Iceberg table (delete, update, insert). Monitor DMS latency and alarm on task errors; validate counts periodically.
Your AWS bill for the data platform jumped 40% last month. How do you investigate?
Open Cost Explorer grouped by service, then by usage type and by cost allocation tag to find the pipeline or resource. Typical culprits are Athena scans on unpartitioned data, crawlers running too often, an EMR cluster left running, NAT gateway data processing, or a new small-files problem. Check Cost Anomaly Detection findings, fix the cause, and add a budget alert so it is caught earlier next time.
Key takeaways
- Run data services in private subnets across two AZs; use an S3 gateway endpoint instead of routing lake traffic through NAT.
- Extract from RDS and Aurora read replicas, snapshot exports, DMS CDC or zero-ETL, never with heavy queries on the primary.
- DynamoDB is designed around access patterns; export to S3 for analytics and use conditional writes for idempotency.
- SQS buffers (visibility timeout, DLQ, at least once), SNS broadcasts, and EventBridge routes by content; SQS and SNS now take 1 MiB messages.
- DMS full load plus CDC lands I/U/D change records; apply the latest change per key with a merge.
- Control cost with activated tags, Cost Explorer, Budgets and anomaly detection, and review designs against the six Well-Architected pillars.
Progress is saved in this browser only. No account needed.

