DEA-C01 · Domain 1: Data Ingestion and Transformation · 34% of the exam
Task 1.1: Perform data ingestion
Getting data in: reading streams from Kinesis, MSK, DynamoDB Streams and AWS DMS, batch loads from S3, databases and SaaS apps through AppFlow, schedules and event triggers that start jobs and crawlers, Lambda consumers of a stream, throttling and rate limits, fan-in and fan-out, and pipelines that can be replayed.
Study it
Streaming ingestion: Kinesis Data Streams, Amazon Data Firehose, MSK and DynamoDB Streams
Try each one before opening the answer. Every option is explained, with the AWS documentation page that proves it.
Question 1 · choose 1
An AWS Lambda function processes records from an Amazon Kinesis data stream that has 4 shards. Each batch finishes well within the function timeout, but the stream's IteratorAge metric keeps growing during business hours. Records that share a partition key must still be processed in order, and the team does not want to reshard the stream. What should a data engineer do?
ARaise the function's reserved concurrency so that more instances can run at once
BIncrease MaximumBatchingWindowInSeconds so each invocation receives more records
CTurn on BisectBatchOnFunctionError
DSet the event source mapping's ParallelizationFactor above 1
Show the answer and why
ARaise the function's reserved concurrency so that more instances can run at once
Incorrect
With the default parallelization factor of 1, Lambda polls one batch at a time from each shard. Setting aside more concurrency does not add parallel work on a 4-shard stream.
BIncrease MaximumBatchingWindowInSeconds so each invocation receives more records
Incorrect
The batching window makes Lambda keep reading until the batch is full or the window expires. It adds waiting time; it does not process a shard in parallel.
CTurn on BisectBatchOnFunctionError
Incorrect
Bisecting splits a batch only when the function returns an error, to isolate bad records. These batches succeed, so the setting changes nothing.
DSet the event source mapping's ParallelizationFactor above 1
Correct
ParallelizationFactor (1 by default, up to 10) sets how many batches Lambda processes from one shard at the same time. AWS suggests it when IteratorAge is high, and Lambda still keeps records in order for each partition key.
A growing IteratorAge with healthy batches means the consumer cannot keep up. Without adding shards, the event source mapping can process several batches per shard at once with ParallelizationFactor, and Lambda keeps per-key order.
Five applications read the same Amazon Kinesis data stream by calling the GetRecords API. Since the fourth and fifth applications were added, all of them receive ReadProvisionedThroughputExceeded errors and their read latency has gone up. Each application must get its own read throughput from every shard, and the team does not want to add shards. Which solution meets these requirements?
ARaise the stream's data retention period from 24 hours to 7 days
BHave the producers aggregate records with the Kinesis Producer Library
CMake every application call GetRecords more often with a smaller limit
DRegister each application as an enhanced fan-out consumer
Show the answer and why
ARaise the stream's data retention period from 24 hours to 7 days
Incorrect
Retention sets how long records stay readable. It gives the consumers no extra read throughput.
BHave the producers aggregate records with the Kinesis Producer Library
Incorrect
Aggregation packs several user records into one Kinesis record on the write side. The consumers still share the same read budget per shard.
CMake every application call GetRecords more often with a smaller limit
Incorrect
Each shard supports up to five GetRecords calls per second and 2 MB per second shared by all callers. More frequent calls hit that limit sooner.
DRegister each application as an enhanced fan-out consumer
Correct
An enhanced fan-out consumer receives up to 2 MB per second per shard that is dedicated to it, independent of the other consumers. Shared consumers split one 2 MB per second budget per shard.
Shared-throughput consumers divide each shard's read capacity among themselves, so every new application slows the others. Enhanced fan-out gives each registered consumer its own pipe from every shard.
A company runs an on-premises MySQL database for its order-entry application. A data engineer must copy the existing 2 TB of data to an Amazon S3 data lake and then keep the lake current with the inserts, updates and deletes made on the source, with little load on the source and no custom code. What should the data engineer do?
ASchedule a nightly AWS Glue job that reads every table over JDBC with job bookmarks turned on
BRun an AWS DMS task that performs a full load and then replicates ongoing changes
CUse AWS DataSync to copy the database's data files to Amazon S3 on an hourly schedule
DExport the tables to CSV each night with a cron job and upload the files with the AWS CLI
Show the answer and why
ASchedule a nightly AWS Glue job that reads every table over JDBC with job bookmarks turned on
Incorrect
For JDBC sources, job bookmarks use key columns to tell new rows from rows already processed. A nightly batch read is not a feed of each update and delete as it happens.
BRun an AWS DMS task that performs a full load and then replicates ongoing changes
Correct
A full-load-and-CDC task copies the existing data and then applies the changes it reads from the source's logs. For MySQL, change capture needs binary logging with binlog_format set to ROW.
CUse AWS DataSync to copy the database's data files to Amazon S3 on an hourly schedule
Incorrect
DataSync moves files and objects between storage systems. Copies of a running database's files are not change records that a data lake can apply.
DExport the tables to CSV each night with a cron job and upload the files with the AWS CLI
Incorrect
This is custom scripting, it reloads whole tables, and it does not show which rows were deleted. Ongoing replication in AWS DMS does this without code.
Change data capture reads the database's own change log, so it adds little load and sees deletes. AWS DMS offers it as a task type that starts with a full load and continues with ongoing replication.
A sales team wants Salesforce opportunity records copied to Amazon S3 every hour. Each run should copy only the records that were created or changed since the previous run, and the company does not want to write or maintain code for the transfer. Which solution meets these requirements?
AAn Amazon AppFlow flow that runs on a schedule with full transfer
BAn Amazon AppFlow flow that runs on demand, started by an operator at the top of each hour
CAn AWS Lambda function on an EventBridge schedule that pages through the Salesforce API
DAn Amazon AppFlow flow that runs on a schedule with incremental transfer
Show the answer and why
AAn Amazon AppFlow flow that runs on a schedule with full transfer
Incorrect
Full transfer copies a snapshot of all records on every run, not only the new and changed ones.
BAn Amazon AppFlow flow that runs on demand, started by an operator at the top of each hour
Incorrect
On-demand flows run when a user starts them. Incremental transfer is an option of schedule-triggered flows, and a person starting runs every hour is not automation.
CAn AWS Lambda function on an EventBridge schedule that pages through the Salesforce API
Incorrect
This works but is custom code that the team would have to write and maintain, which the requirement rules out.
DAn Amazon AppFlow flow that runs on a schedule with incremental transfer
Correct
Schedule-triggered flows can use incremental transfer, which copies only the records created or changed since the last run. AppFlow connects to Salesforce without code.
AppFlow is the no-code path for SaaS sources such as Salesforce. The trigger (on demand, on event or on schedule) and the transfer mode (full or incremental) decide what each run copies.
A data engineer must load 600 GB of CSV files from Amazon S3 into a DynamoDB table that a new application will start using next week. A trial load with an AWS Glue job was throttled and needed a large amount of provisioned write capacity. What is the MOST cost-effective way to do the initial load?
ASwitch the table to on-demand capacity mode and run the same AWS Glue job again
BUse DynamoDB import from S3 to create a new table from the files
CWrite the rows with BatchWriteItem calls and retry throttled items with backoff
DGive the AWS Glue job more workers
Show the answer and why
ASwitch the table to on-demand capacity mode and run the same AWS Glue job again
Incorrect
On-demand mode charges for every write the job performs. It removes the capacity planning, not the cost of 600 GB of individual writes.
BUse DynamoDB import from S3 to create a new table from the files
Correct
Import from S3 creates a new table from CSV, DynamoDB JSON or Amazon Ion files and does not consume write capacity, so no extra capacity has to be provisioned for the load.
CWrite the rows with BatchWriteItem calls and retry throttled items with backoff
Incorrect
Batching and backoff spread the same writes over more time. Each item still consumes write capacity on the table.
DGive the AWS Glue job more workers
Incorrect
More workers send writes faster, so they hit the table's write capacity sooner. The capacity the load needs stays the same.
Writing through the table's API pays for every item. Import from S3 is priced on the source data size and leaves the table's write capacity alone, which fits a one-time initial load into a new table.
An application writes events to an Amazon Kinesis data stream that keeps the default retention period. Consumer teams sometimes deploy a bug and notice it up to 3 days later. After a fix, the consumer must reprocess the events that arrived from the moment the bug was deployed, but not older events. Which actions should a data engineer take? (Choose TWO.)
AIncrease the stream's retention period to 7 days
BRestart the fixed consumer with a TRIM_HORIZON shard iterator on every shard
CRestart the fixed consumer with an AT_TIMESTAMP shard iterator set to the deployment time
DRegister the consumer as an enhanced fan-out consumer of the stream
ESwitch the stream from provisioned mode to on-demand capacity mode
Show the answer and why
AIncrease the stream's retention period to 7 days
Correct
A stream keeps records for 24 hours by default and up to 8,760 hours (365 days). Seven days covers a bug found 3 days late.
BRestart the fixed consumer with a TRIM_HORIZON shard iterator on every shard
Incorrect
TRIM_HORIZON starts at the oldest record still in the shard, which would also reprocess the events from before the bug.
CRestart the fixed consumer with an AT_TIMESTAMP shard iterator set to the deployment time
Correct
AT_TIMESTAMP starts reading at a given point in time, so the consumer replays exactly the events since the bad deployment.
DRegister the consumer as an enhanced fan-out consumer of the stream
Incorrect
Enhanced fan-out gives a consumer dedicated read throughput. It does not keep records any longer.
ESwitch the stream from provisioned mode to on-demand capacity mode
Incorrect
The capacity mode decides how shards and throughput are managed and billed. Retention is a separate setting.
Replay needs two things: the records must still be in the stream, and the consumer must be able to start reading at a chosen point. A longer retention period gives the first, an AT_TIMESTAMP iterator the second.
Sensors write about 4,000 readings per second to an Amazon Kinesis data stream that has 2 provisioned shards. Each reading is about 200 bytes. The producers are throttled even though each shard receives far less than 1 MB per second. Which change removes the throttling with the smallest increase in cost?
AGive every reading a random partition key so that both shards share the load
BUse the Kinesis Producer Library and aggregate many readings into each record
CRegister the downstream applications as enhanced fan-out consumers of the stream
DIncrease the stream's data retention period from 24 hours to 7 days
Show the answer and why
AGive every reading a random partition key so that both shards share the load
Incorrect
Spreading the readings evenly still sends about 2,000 records per second to each shard, twice the per-shard record limit.
BUse the Kinesis Producer Library and aggregate many readings into each record
Correct
Each shard accepts up to 1,000 records per second, a limit that binds when records are small. Aggregation stores many readings in one Kinesis record, so the record rate drops without adding shards.
CRegister the downstream applications as enhanced fan-out consumers of the stream
Incorrect
Enhanced fan-out gives consumers their own read throughput. It does not change how many records producers can write to a shard.
DIncrease the stream's data retention period from 24 hours to 7 days
Incorrect
The retention period controls how long records stay readable in the stream, not how fast producers can write them.
The bottleneck is the record count, not the bytes: 2,000 small records per shard per second against a limit of 1,000. KPL aggregation packs the readings into fewer, larger Kinesis records and fixes it without new shards.
An Amazon Data Firehose stream receives about 200 KB of JSON per second and delivers it to Amazon S3 with the default buffering settings. The objects in S3 are about 5 MB each. The analytics team wants objects close to 128 MB and accepts up to 15 minutes of delivery delay. Which TWO changes should a data engineer make? (Choose TWO.)
ASet the buffer size to 128 MiB
BSet the buffer interval to 60 seconds
CSet the buffer interval to 900 seconds
DTurn on dynamic partitioning by customer ID
EAdd a Lambda function that transforms the incoming records
Show the answer and why
ASet the buffer size to 128 MiB
Correct
Firehose delivers when the first buffer condition is met. A 128 MiB size hint lets each object grow to the size the team wants.
BSet the buffer interval to 60 seconds
Incorrect
A shorter interval is met sooner, so Firehose would deliver even smaller objects.
CSet the buffer interval to 900 seconds
Correct
At about 200 KB per second, 128 MiB takes more than 10 minutes to arrive. The default 300-second interval would flush first, so the interval must rise to its maximum of 900 seconds.
DTurn on dynamic partitioning by customer ID
Incorrect
Dynamic partitioning groups records by a key into separate S3 prefixes. It splits the data further instead of making each object larger.
EAdd a Lambda function that transforms the incoming records
Incorrect
A transformation function changes the records before delivery. The buffering hints still decide when Firehose writes an object.
Size and interval work together: Firehose flushes on whichever is reached first. To get 128 MB objects from a slow stream, raise both the size hint and the interval, accepting up to 15 minutes of delay.
A partner publishes order changes through a REST API that returns every change since a given timestamp. No managed connector exists for the partner. A data engineer must call the API every 10 minutes and write the response to Amazon S3. Each call finishes in under a minute, and the team does not want to manage servers. Which solution meets these requirements?
AAn Amazon S3 event notification on the bucket that invokes a Lambda function
BAn Amazon EventBridge Scheduler schedule that invokes an AWS Lambda function
CAn AWS DMS task that uses the partner API as its source endpoint
DAn Amazon Data Firehose stream that uses the partner API as its source
Show the answer and why
AAn Amazon S3 event notification on the bucket that invokes a Lambda function
Incorrect
S3 event notifications fire on events in the bucket, such as new objects. Nothing in the bucket changes on a 10-minute clock.
BAn Amazon EventBridge Scheduler schedule that invokes an AWS Lambda function
Correct
EventBridge Scheduler runs rate or cron schedules and can invoke Lambda as a target, so a short function can call the API and write the result to S3 without servers.
CAn AWS DMS task that uses the partner API as its source endpoint
Incorrect
AWS DMS replicates from database engines and a few AWS data sources, not from a REST API.
DAn Amazon Data Firehose stream that uses the partner API as its source
Incorrect
A Firehose stream receives data that producers write to it or reads a Kinesis data stream or Amazon MSK. It does not poll an API.
Pulling from an API on a timer is a schedule plus code: EventBridge Scheduler for the timer and a Lambda function for the call, both serverless.
An AWS Glue job uses a network connection in a VPC so that it can read from an Amazon RDS database. The job must also call a partner's HTTPS endpoint, and the partner allows traffic only from one fixed public IP address. What should a data engineer do?
ARoute the connection's private subnet through a public NAT gateway and share its Elastic IP address
BMove the connection to a public subnet with an internet gateway route and share the job's public IP
CSend the traffic through a private NAT gateway and give the partner that gateway's IP address
DTurn on public IP assignment in the job properties and allow-list each worker address
Show the answer and why
ARoute the connection's private subnet through a public NAT gateway and share its Elastic IP address
Correct
Glue gives its network interfaces only private addresses, so internet access needs a NAT gateway. A public NAT gateway has an Elastic IP address, which the partner can allow.
BMove the connection to a public subnet with an internet gateway route and share the job's public IP
Incorrect
Glue assigns no public IP addresses to the network interfaces it creates, so a route to an internet gateway alone gives the job no public address to share.
CSend the traffic through a private NAT gateway and give the partner that gateway's IP address
Incorrect
A private NAT gateway cannot have an Elastic IP address and is meant for other VPCs or on-premises networks, not the internet.
DTurn on public IP assignment in the job properties and allow-list each worker address
Incorrect
Glue gives its network interfaces only private addresses from the subnet, so the workers have no public addresses to list.
Glue jobs in a VPC reach the internet through NAT. A public NAT gateway gives all of that traffic one stable Elastic IP address for the partner's allow list.
Every insert, update, and delete in an on-premises PostgreSQL database must reach three independent applications within seconds. Each application reads the changes at its own pace. Which solution meets these requirements?
AAn AWS DataSync task that copies the database's data files to Amazon S3 every minute
BAn AWS DMS task with ongoing replication that targets an Amazon Kinesis data stream
CAn Amazon AppFlow flow that uses the PostgreSQL database as its source
DAn hourly AWS Glue job with job bookmarks that writes new rows to Amazon S3
Show the answer and why
AAn AWS DataSync task that copies the database's data files to Amazon S3 every minute
Incorrect
DataSync transfers file and object data between storage systems. It does not capture row-level database changes.
BAn AWS DMS task with ongoing replication that targets an Amazon Kinesis data stream
Correct
AWS DMS can capture ongoing changes and publish them as JSON records to a Kinesis data stream, which several applications can read in real time.
CAn Amazon AppFlow flow that uses the PostgreSQL database as its source
Incorrect
AppFlow exchanges data between SaaS applications and AWS services, not from a self-managed PostgreSQL database.
DAn hourly AWS Glue job with job bookmarks that writes new rows to Amazon S3
Incorrect
Job bookmarks let a scheduled job skip data it has already processed, but an hourly batch cannot deliver changes within seconds.
Change data capture into a stream gives many consumers the same ordered feed. AWS DMS with a Kinesis Data Streams target does both parts.