AWS courseLesson 10 of 12
AWS course · Lesson 10 of 12
Monitoring Data Pipelines on AWS: CloudWatch and EventBridge
Monitor AWS data pipelines with CloudWatch metrics, alarms, custom metrics, Logs Insights and dashboards, plus EventBridge rules that react to job failures.
On this page
A pipeline that fails silently is worse than one that fails loudly: dashboards show yesterday’s numbers and nobody knows. On AWS, Amazon CloudWatch collects metrics and logs from every service and from your own code, and Amazon EventBridge reacts to state changes such as “Glue job failed”. This lesson shows how to use both so that you hear about a late, failed or wrong pipeline before your users do.
Sample code
The runnable examples are simplified models: how an alarm evaluates datapoints, what an Embedded Metric Format log line looks like, and how an EventBridge pattern matches events. The CLI commands and queries need an AWS account and are not executed.
Metrics and alarms
What it is. A metric is a time series of datapoints identified by a namespace (for example AWS/Lambda, AWS/Kinesis, Glue), a name and dimensions (key-value pairs such as FunctionName=validate-file). An alarm watches one metric (or a metric math expression) and changes state when it crosses a threshold.
How it works.
- Datapoints are aggregated over a period (for example 60 or 300 seconds) with a statistic:
Sum,Average,Minimum,Maximum,SampleCountor a percentile such asp95. - An alarm has three states:
OK,ALARMandINSUFFICIENT_DATA. It evaluates the last N periods (evaluation periods) and goes toALARMwhen M of them breach (datapoints to alarm). “3 out of 5” ignores a single spike but catches a sustained problem. - Missing data is treated as
missing,notBreaching,breachingorignore. For a pipeline that should report regularly, silence is itself a failure, sobreachingis often right. - Alarm actions notify an SNS topic (email, chat, paging tools), or trigger Auto Scaling or EC2 actions. Composite alarms combine several alarms with
ANDandORto cut noise, and anomaly detection alarms learn a metric’s normal band.
This model evaluates a freshness metric (minutes since the last successful load) with a “3 out of 5 periods above 60 minutes” alarm:
def alarm_state(datapoints, threshold, m, n, treat_missing="missing"):
"""CloudWatch 'M out of N' evaluation for a GreaterThanThreshold alarm.
datapoints: newest last; None means no data for that period."""
window = datapoints[-n:]
if treat_missing == "breaching":
window = [threshold + 1 if d is None else d for d in window]
elif treat_missing == "notBreaching":
window = [threshold if d is None else d for d in window]
present = [d for d in window if d is not None]
if not present:
return "INSUFFICIENT_DATA"
breaching = sum(d > threshold for d in present)
return "ALARM" if breaching >= m else "OK"
# Minutes since the last successful load ("freshness"), one datapoint per 5-minute period.
spiky = [20, 95, 30, 25, 20]
late = [70, 85, 100, 115, 130]
silent = [20, None, None, None, None]
for name, series in [("one spike", spiky), ("really late", late), ("job stopped reporting", silent)]:
print(f"{name:<22} 3 of 5 > 60 min -> {alarm_state(series, 60, 3, 5):<17}"
f" (missing=breaching -> {alarm_state(series, 60, 3, 5, 'breaching')})")
one spike 3 of 5 > 60 min -> OK (missing=breaching -> OK)
really late 3 of 5 > 60 min -> ALARM (missing=breaching -> ALARM)
job stopped reporting 3 of 5 > 60 min -> OK (missing=breaching -> ALARM)
The third case is the important one: a job that stops reporting looks healthy unless missing data is treated as breaching.
Useful built-in metrics for data pipelines:
| Service | Metric | Watch for |
|---|---|---|
| Lambda | Errors, Throttles, Duration, IteratorAge (stream sources), ConcurrentExecutions |
Failures, throttling, consumers falling behind |
| Kinesis Data Streams | GetRecords.IteratorAgeMilliseconds, WriteProvisionedThroughputExceeded, ReadProvisionedThroughputExceeded |
Lag and hot shards |
| Amazon Data Firehose | DeliveryToS3.DataFreshness, DeliveryToS3.Success |
Data not landing |
| SQS | ApproximateAgeOfOldestMessage, ApproximateNumberOfMessagesVisible (also on DLQs) |
Backlog, poison messages |
| Step Functions | ExecutionsFailed, ExecutionsTimedOut, ExecutionTime |
Pipeline failures and slowdowns |
| Glue | Job metrics in the Glue namespace (when enabled), plus job state events |
Slow stages, failed tasks |
| Athena | ProcessedBytes per workgroup |
Scan cost spikes |
aws cloudwatch put-metric-alarm --alarm-name orders-freshness \
--namespace DataPipelines --metric-name MinutesSinceLastLoad \
--dimensions Name=Pipeline,Value=orders \
--statistic Maximum --period 300 \
--evaluation-periods 5 --datapoints-to-alarm 3 \
--threshold 60 --comparison-operator GreaterThanThreshold \
--treat-missing-data breaching \
--alarm-actions arn:aws:sns:eu-west-1:111122223333:pipeline-alerts
Pitfalls.
- Alarms on
Averagehide spikes; useMaximumor a percentile for latency and lag. - One datapoint out of one: alarms that flap on every blip train people to ignore them.
- Alarms nobody owns. Every alarm should have an owner and a runbook link in its description.
In interviews. Explain namespace, dimensions, period and statistic, “M out of N” evaluation and missing-data treatment, and give two or three metrics you would alarm on for a specific pipeline.
Custom metrics
What it is. Metrics your code publishes about things only it knows: rows read and written, rows rejected by validation, minutes since the last successful load, bytes per partition.
How it works. Two ways to publish:
PutMetricDataAPI calls (CLI or SDK). Simple, but each call is a request, so batch values and avoid calling it per record.- Embedded Metric Format (EMF): write a structured JSON log line to CloudWatch Logs (from Lambda, a container or any agent), and CloudWatch extracts the metrics asynchronously. Other fields stay in the log for Logs Insights, which links metrics and context.
aws cloudwatch put-metric-data --namespace DataPipelines \
--metric-name RowsWritten --unit Count --value 120431 \
--dimensions Pipeline=orders,Stage=curated
import json, time
def emf_line(pipeline, stage, rows_written, rejected, duration_ms):
"""Embedded Metric Format: a structured log line that CloudWatch turns into custom metrics."""
return json.dumps({
"_aws": {
"Timestamp": int(time.time() * 1000),
"CloudWatchMetrics": [{
"Namespace": "DataPipelines",
"Dimensions": [["Pipeline", "Stage"]],
"Metrics": [
{"Name": "RowsWritten", "Unit": "Count"},
{"Name": "RowsRejected", "Unit": "Count"},
{"Name": "DurationMs", "Unit": "Milliseconds"},
],
}],
},
"Pipeline": pipeline, "Stage": stage,
"RowsWritten": rows_written, "RowsRejected": rejected, "DurationMs": duration_ms,
"run_id": "2026-10-05T02:30", # extra keys stay searchable in Logs Insights
})
line = json.loads(emf_line("orders", "curated", 120_431, 12, 482_000))
print(sorted(k for k in line if k != "_aws"))
print([m["Name"] for m in line["_aws"]["CloudWatchMetrics"][0]["Metrics"]])
['DurationMs', 'Pipeline', 'RowsRejected', 'RowsWritten', 'Stage', 'run_id']
['RowsWritten', 'RowsRejected', 'DurationMs']
Pitfalls.
- Dimension cardinality. Each unique combination of dimension values is a separate metric, billed separately. Never use run IDs, file names or user IDs as dimensions; put them in log fields instead.
- Publishing only success counts: a job that crashes before publishing reports nothing. Use a freshness metric with missing data treated as breaching, or alarm on the absence of success events.
- High-resolution (one-second) metrics cost more and are rarely needed for batch pipelines.
In interviews. The best custom metrics for pipelines are freshness, volume (rows in and out), and rejects. Mention EMF and the cardinality trap.
Logs and Logs Insights
What it is. CloudWatch Logs stores log events in log groups (one per Lambda function, Glue job type, Step Functions state machine and so on) made of log streams. Logs Insights is a query language for searching and aggregating them.
How it works.
- Set a retention period on every log group; the default is to keep logs forever.
- Metric filters turn matching log lines into metrics (for example, count lines containing
ERROR). - Subscription filters stream logs to Lambda, Kinesis or Firehose for processing or archiving in S3.
- Logs Insights queries pipe commands:
fields,filter,parse,stats,sort,limit. LambdaREPORTlines expose fields such as@durationand@maxMemoryUsed.
fields @timestamp, @message
| filter @message like /ERROR|Exception/
| stats count(*) as errors by bin(15m)
| sort errors desc
filter @type = "REPORT"
| stats avg(@duration) as avg_ms, pct(@duration, 95) as p95_ms,
max(@maxMemoryUsed / 1000 / 1000) as max_mem_mb by bin(1h)
filter Pipeline = "orders" and RowsRejected > 0
| fields @timestamp, run_id, RowsWritten, RowsRejected
| sort @timestamp desc
| limit 20
The third query reads the fields of the EMF log lines from the previous section, which is why structured JSON logs are worth the effort.
Pitfalls.
- Unstructured print statements that are hard to filter. Log JSON with consistent keys (pipeline, run ID, stage, counts).
- Logs Insights bills by data scanned; narrow the time range and log groups.
- No retention set, so log storage grows forever.
In interviews. Write a Logs Insights query from memory (filter, stats by bin), and explain structured logging and retention.
Dashboards
What it is. A CloudWatch dashboard is a page of widgets (metric graphs, single numbers, alarm status, Logs Insights results, text) shared by a team.
How it works. Build dashboards in the console or as JSON (ideal for infrastructure as code). A good pipeline dashboard answers, at a glance: did today’s runs succeed, is data fresh, are volumes normal, is anything backing up, and what is it costing?
{
"widgets": [
{
"type": "metric",
"x": 0, "y": 0, "width": 12, "height": 6,
"properties": {
"title": "Rows written per hour",
"region": "eu-west-1",
"stat": "Sum",
"period": 3600,
"metrics": [["DataPipelines", "RowsWritten", "Pipeline", "orders", "Stage", "curated"]]
}
},
{
"type": "alarm",
"x": 12, "y": 0, "width": 12, "height": 6,
"properties": {
"title": "Pipeline alarms",
"alarms": ["arn:aws:cloudwatch:eu-west-1:111122223333:alarm:orders-freshness"]
}
}
]
}
aws cloudwatch put-dashboard --dashboard-name data-pipelines \
--dashboard-body file://dashboard.json
Pitfalls. Dashboards are for humans looking; they do not wake anyone up. Every important condition on a dashboard also needs an alarm.
In interviews. Describe the panels you would put on a pipeline dashboard (status, freshness, volume, backlog, errors, cost) and keep dashboards in code.
EventBridge rules
What it is. Amazon EventBridge is an event bus. AWS services publish events about state changes (a Glue job failed, a Step Functions execution finished, an object was created in S3 when the bucket sends events to EventBridge). Rules match events with event patterns and send them to targets such as SNS, SQS, Lambda, Step Functions or another bus. EventBridge also runs schedules, and EventBridge Scheduler is the recommended service for cron and rate schedules.
How it works. An event pattern lists the values a field must have; fields not mentioned are ignored, and arrays mean “any of”. Content filters such as prefix, anything-but, numeric and exists give more control. Glue publishes Glue Job State Change events for SUCCEEDED, FAILED, TIMEOUT and STOPPED. This matcher implements a small subset of the pattern rules to show which events a failure rule catches:
def matches(pattern, event):
"""A small subset of EventBridge pattern matching: exact values, prefix and anything-but."""
for key, rule in pattern.items():
if key not in event:
return False
value = event[key]
if isinstance(rule, dict):
if not matches(rule, value):
return False
continue
ok = False
for cond in rule:
if isinstance(cond, dict) and "prefix" in cond:
ok |= isinstance(value, str) and value.startswith(cond["prefix"])
elif isinstance(cond, dict) and "anything-but" in cond:
ok |= value not in cond["anything-but"]
else:
ok |= value == cond
if not ok:
return False
return True
pattern = json.loads("""
{
"source": ["aws.glue"],
"detail-type": ["Glue Job State Change"],
"detail": {
"jobName": [{ "prefix": "orders-" }],
"state": ["FAILED", "TIMEOUT"]
}
}
""")
events = [
{"source": "aws.glue", "detail-type": "Glue Job State Change",
"detail": {"jobName": "orders-curated", "state": "FAILED", "jobRunId": "jr_1"}},
{"source": "aws.glue", "detail-type": "Glue Job State Change",
"detail": {"jobName": "orders-curated", "state": "SUCCEEDED", "jobRunId": "jr_2"}},
{"source": "aws.glue", "detail-type": "Glue Job State Change",
"detail": {"jobName": "clicks-hourly", "state": "TIMEOUT", "jobRunId": "jr_3"}},
]
for e in events:
print(e["detail"]["jobName"], e["detail"]["state"], "->", "alert" if matches(pattern, e) else "ignored")
orders-curated FAILED -> alert
orders-curated SUCCEEDED -> ignored
clicks-hourly TIMEOUT -> ignored
aws events put-rule --name glue-orders-failures --event-pattern file://pattern.json
aws events put-targets --rule glue-orders-failures \
--targets 'Id=alerts,Arn=arn:aws:sns:eu-west-1:111122223333:pipeline-alerts'
An input transformer can turn the raw event into a readable message (“Glue job orders-curated FAILED, run jr_1”) before it reaches SNS. Targets have a retry policy and can have a dead-letter queue for events that cannot be delivered.
Pitfalls.
- Patterns that are too broad (all of
aws.glue) flood the channel; too narrow (a typo in the job name) catch nothing. Test patterns with the console’s sandbox ortest-event-pattern. - Forgetting the target’s resource policy (SNS topic policy, Lambda permission) so EventBridge cannot deliver.
- Using events alone for SLAs. An event tells you something happened; a freshness alarm tells you something did not happen.
In interviews. Show a pattern for failed Glue jobs routed to SNS, mention input transformers and DLQs, and contrast event-driven alerts (failures) with metric alarms (lateness and drift).
Monitoring pipelines
Pulling it together, a monitored AWS pipeline watches four things:
| Signal | Question | How on AWS |
|---|---|---|
| Failures | Did any step fail? | EventBridge rules on Glue, Step Functions, EMR and Batch state changes; alarms on Lambda Errors and DLQ depth |
| Freshness | Is the newest data recent enough? | Custom MinutesSinceLastLoad metric (or Firehose DataFreshness), alarm with missing data as breaching |
| Volume and quality | Did we load roughly the expected rows, with few rejects? | Custom row-count metrics, Glue Data Quality results and events, anomaly detection alarms |
| Performance and cost | Is it getting slower or more expensive? | Duration metrics, Kinesis and Lambda iterator age, Athena ProcessedBytes, AWS Budgets and Cost Anomaly Detection |
A practical setup:
- Every job emits structured logs with a run ID, and EMF metrics for rows in, rows out and rejects.
- Step Functions (or Airflow) catches failures and publishes to an SNS alerts topic; an EventBridge rule covers anything started outside the orchestrator.
- A freshness alarm per important table, with missing data treated as breaching.
- DLQ depth alarms on every queue and Lambda destination.
- One dashboard per pipeline, defined in code, and runbooks linked from alarm descriptions.
- Cost monitoring through tags, Budgets and anomaly detection (covered in the landscape lesson).
Pitfalls. Alerting on everything. Page only on what needs action now (a failed critical load, stale data near an SLA); send the rest to a channel or a ticket.
In interviews. “How would you know this pipeline broke?” Answer with all four signals, not just “CloudWatch alarms”, and explain why freshness catches failures that error alerts miss.
Practice questions
A daily load silently stopped running last week and nobody noticed. What monitoring would have caught it?
A freshness metric (minutes since last successful load, or the max partition date) with an alarm that treats missing data as breaching. Failure events would not fire because nothing ran. An EventBridge rule on scheduler or orchestrator failures, and a row-count anomaly alarm, would add coverage.
Why should you not use a run ID as a CloudWatch metric dimension?
Every unique set of dimension values creates a separate metric, billed separately and impossible to alarm on as a series. Run IDs are unbounded, so they explode the number of metrics. Put run IDs in log fields (for example in EMF log lines) and keep dimensions low-cardinality, such as pipeline and stage.
Write an EventBridge pattern that matches failed or timed-out Glue jobs whose names start with “orders-”.
{"source": ["aws.glue"], "detail-type": ["Glue Job State Change"], "detail": {"jobName": [{"prefix": "orders-"}], "state": ["FAILED", "TIMEOUT"]}}, with an SNS topic (through an input transformer) as the target.
An alarm on a Kinesis consumer’s iterator age keeps flapping. How do you fix it?
Use “M out of N” evaluation (for example 3 of 5 periods) so short spikes do not alarm, use the Maximum statistic over a sensible period, set a threshold tied to the SLA rather than to normal jitter, and consider a composite alarm that also requires errors or throttles. Then find the cause of the spikes (hot shard, slow batches).
Write a Logs Insights query for the 95th percentile duration of a Lambda function per hour.
filter @type = "REPORT" | stats pct(@duration, 95) as p95_ms by bin(1h), run against the function’s log group for the time range of interest.
Key takeaways
- CloudWatch metrics are namespaced, dimensioned time series; alarms evaluate M out of N periods with a chosen statistic and missing-data rule.
- Publish custom pipeline metrics (freshness, volume, rejects) with PutMetricData or EMF, and keep dimensions low-cardinality.
- Log structured JSON, set log retention, and query with Logs Insights.
- Dashboards are for people; alarms and EventBridge rules are what notify them.
- EventBridge rules match service state changes, such as failed Glue jobs, and route them to SNS, SQS, Lambda or Step Functions.
- Monitor failures, freshness, volume and quality, and performance and cost; freshness alarms catch the silent failures.
Progress is saved in this browser only. No account needed.