
Monitoring a Pipeline with Estuary’s Observability Tools
Learn how to monitor a real-time data pipeline with Estuary’s native dashboard, alerts, flowctl, and OpenMetrics API, then extend monitoring with Prometheus to track freshness, throughput, backlog, errors, and overall pipeline health.

Unlike traditional monitoring, which is often reactive and focused on predefined system metrics, data observability provides broader visibility into how data moves through a pipeline. For real-time data pipelines, this means monitoring signals such as data freshness, throughput, processing lag, errors, and task health so teams can quickly detect when data is delayed, stalled, or failing downstream.
Data observability is the continuous monitoring of the health, reliability, and performance of data as it moves from ingestion through processing and downstream delivery. By tracking signals such as freshness, throughput, backlog, errors, and task status, teams can better understand whether data is arriving and being delivered reliably and as expected.
In this tutorial, we’ll build a real-time data pipeline and explore Estuary’s native observability capabilities through its dashboard, built-in alerts, logs, and flowctl, Estuary’s command-line interface. We’ll then use Estuary’s OpenMetrics API with Prometheus for more granular monitoring, querying, visualization, and custom alerting. Along the way, we’ll track pipeline freshness, throughput, backlog, and errors to identify whether an issue originates at the source, within the pipeline, or at the destination.
What Should You Monitor in a Data Pipeline?
A healthy data pipeline is not just one that is running. To understand whether data is moving reliably from source to destination, teams should monitor a small set of operational signals: freshness, throughput, backlog, errors, and task health.
| Signal | What it tells you | Example Estuary metric or signal |
|---|---|---|
| Freshness | How recently new data was processed or delivered | captured_last_published_at_time_seconds |
| Throughput | How quickly records are moving through the pipeline | rate(captured_in_docs_total[5m]) |
| Backlog | Whether downstream processing is keeping up with incoming data | materialized_bytes_behind |
| Errors | Whether tasks are producing runtime errors | logged_errors_total |
| Task health | Whether a task has failed, stalled, become idle, or stopped making progress | Estuary native alerts |
These signals are most useful when monitored together. For example, healthy capture throughput combined with a growing materialization backlog can indicate that the destination is falling behind, while zero capture throughput may point to an upstream ingestion problem.
Throughout this tutorial, we’ll use Estuary’s dashboard, native alerts, flowctl, and OpenMetrics metrics in Prometheus to monitor these signals across the pipeline.
Architecture Overview
The architecture separates the data path from the observability layer. The data path shows how records move through the pipeline, while the observability layer shows how pipeline activity is inspected, measured, and monitored without changing the underlying flow of data.
In this demo, a Python producer generates synthetic transaction events and sends them to an Estuary HTTP Webhook capture. Estuary writes the incoming records to a collection, which is then materialized to Google Sheets as the downstream destination. Throughout this process, Estuary collects operational telemetry that can be inspected through its native dashboard, alerts, and flowctl, or exported through the OpenMetrics API for external monitoring with Prometheus.
Steps Taken
- Generate synthetic data — A Python producer creates transaction events at a controlled rate to simulate a continuously updating source.
- Ingest data through the HTTP Webhook — The producer sends each record to Estuary using HTTP POST requests, where the Webhook capture receives and processes the events.
- Store the captured events in an Estuary collection.
- Access captures & metrics through multiple interfaces
- Dashboard: Browse collections, provides task status, document and byte counts, logs, and alerts
- flowctl: Runtime statistics, logs, and alert configuration.
- OpenMetrics API: Estuary makes capture, derivation, materialization, and log metrics available in a Prometheus-compatible format.
- Collect and query metrics with Prometheus — Prometheus periodically scrapes the OpenMetrics endpoint and stores the metrics as time-series data for analyzing freshness, throughput, backlog, and errors.
- Materialize into Google Sheets & track metrics: Configure Estuary to materialize your captured collections into Google Sheets as a destination.
Let’s go ahead and build this pipeline!
Build the Real-Time Data Pipeline
Prerequisites
Before getting started, make sure you have the following:
- An Estuary account: You’ll need access to an Estuary account where you can create a capture and materialization.
- A Python environment: A virtual environment such as venv is recommended for installing and managing the script’s dependencies.
- Estuary flowctl: The Estuary command-line tool is useful for interacting with collections, captures, and pipeline metadata programmatically.
- A Prometheus instance: Required to scrape and visualize metrics exposed through Estuary’s OpenMetrics API. This tutorial uses Prometheus for metric collection and monitoring.
- A Google Sheets: Required if you want to follow the Google Sheets materialization portion of the tutorial.
Step 1: Generate Synthetic Data With Python
First, we will generate realistic synthetic transaction data using a Python script that uses Faker to generate synthetic transaction data and sends it to your Estuary webhook endpoint via HTTP POST at regular intervals. This simulates continuous real-world data ingestion.
webhook_producer.py:
plaintextimport random
import time
import requests
from datetime import datetime, timezone
from faker import Faker
fake = Faker()
WEBHOOK_URL = "https://your-webhook-url/sample_data"
BEARER_TOKEN = "your_token_here"
RECORDS_PER_SECOND = 2
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {BEARER_TOKEN}"
}
def generate_transaction():
return {
"transaction_id": fake.uuid4(),
"customer_name": fake.name(),
"product": random.choice(["Laptop", "Phone", "Tablet", "Headphones"]),
"amount": round(random.uniform(20, 1500), 2),
"timestamp": datetime.now(timezone.utc).isoformat()
}
print("Sending transactions to Estuary. Press Ctrl+C to stop.")
try:
while True:
transaction = generate_transaction()
response = requests.post(
WEBHOOK_URL,
json=transaction,
headers=headers,
timeout=10
)
print(
f"{response.status_code} | "
f"{transaction['product']} | "
f"${transaction['amount']}"
)
time.sleep(1 / RECORDS_PER_SECOND)
except KeyboardInterrupt:
print("\nProducer stopped.")
This script generates synthetic transaction records with Faker and continuously sends them to the Estuary HTTP Webhook at a configurable rate. It runs until Ctrl+C is pressed, providing a steady stream of events for monitoring capture activity, throughput, freshness, and downstream materialization behavior.
Step 2: Set up an Estuary source connector for HTTP webhooks
- Open your browser and go to https://dashboard.estuary.dev
- Log in with your Estuary account credentials
- Click Create Capture (or navigate to Captures → New Capture)
- In the connector search box, type "Webhook"
- Output Collection: Enter a collection name where data will be stored. Example: RuheeShrestha/sample_data/source-http-ingest
- Click Publish to create and start the connector. After publishing, Estuary displays your unique webhook endpoint
- Find it by clicking on your newly created capture
- Look for the webhook URL (it will look like): **https://5d1503c27edd36da-8080.reactor.aws-us-east-1-c1.dp.estuary-data.com/**
- Copy the entire URL endpoint and add /sample_data and add it to your Python Script.
- Add Bearer token as well optionally if you set it up in the capture.
- Run the Python script and you will see an output with something like this:
And your capture will start ingesting the data:
Step 3: Configure the Estuary HTTP Webhook Capture
Estuary Native Observability
Estuary already provides several ways to inspect pipeline activity and identify operational issues. The dashboard provides an at-a-glance view of task health, native alerts notify you when common failure conditions occur, and flowctl provides more detailed command-line access to logs, statistics, and alert configuration.
Method #1: Estuary Dashboard
You can view a subset of logs and statistics for individual tasks in Estuary’s Dashboard. After you publish a new capture or materialization, in the capture details you can find an at-a-glance health view of: Overview, Alerts, Specs, History and Logs.
- Overview: provides the resulting runtime information, including the capture's endpoint, target collection, shard status, data plane, and number of documents/data written over time.
Two statistics are shown for each capture, collection, and materialization:
- Bytes Written or Read: This corresponds to the bytesTotal property of the stats collection.
- Docs Written or Read: This corresponds to the docsTotal property of the stats collection.
- Alerts: active and resolved operational alerts.
- Specs: The Specs tab shows the configuration that defines the Estuary task. It provides the technical definition of how the capture is configured and how incoming data is written into Estuary collections. Read more here
- History: The History tab records changes made to the capture over time. It lets you review previous versions of the specification and see how the configuration has changed between deployments.
- Logs: operational events, warnings, and errors produced by the task.
Native Alerts
Estuary's native alerting system monitors common operational conditions without requiring an external monitoring stack. Alert subscriptions can be scoped by catalog prefix and routed to specific recipients.
You can set up an alert system where you can define:
- what notifications you want to receive,
- who should receive them via email,
- which environments/prefixes you want those settings to apply to and
- when alerts should fire with fine tuned thresholds.
Available Alert Types
You can create a custom alert subscription using Estuary's built-in email alerts, tracking the following:
| Alert | Fires when |
|---|---|
| Auto-Discovery Failed | Automatic discovery of new or updated source resources fails, such as when Estuary can no longer authenticate with the source. |
| Background Publication Failed | An automated background update to a Flow specification fails to publish. This can affect schema updates, materialization bindings, or other automated spec changes. |
| Data Movement Stalled | A capture or materialization receives no new documents within its configured interval. The interval can be configured per task or inherited from a prefix-level default. |
| Task Failed | A task repeatedly fails even after Estuary attempts automatic recovery. By default, the alert fires after 3 failures within an 8-hour window. |
| Task Chronically Failing | A task continues failing for an extended period and is unable to make progress. |
| Task Auto-Disabled (Failing) | A chronically failing task remains unhealthy long enough that Estuary disables it until the underlying problem is addressed. |
| Task Idle | A task has not processed data for an extended period and has not been modified recently. |
| Task Auto-Disabled (Idle) | A task remains idle long enough that Estuary automatically disables it. |
Setting up your own alerts
To start configuring your notifications:
- Navigate to the Admin Settings page in your Estuary dashboard.
- Click Configure Notifications under the Organization Notifications section. The alert subscription modal will open.
- Choose your prefix: the base tenant name (such as sample_data/) will apply notification settings to the whole tenant. Typing in a subprefix (such as sample_data/source-http-ingest) will apply your settings to tasks under that specific prefix. This is how you can separate default notification settings by environment, or even down to the task level.
- Set desired global settings. These are default values that will apply to the chosen prefix. Currently, you can choose a default setting for Data Movement Stalled alerts.
- Define one or more recipients to receive notifications. Each recipient consists of an email (this can be a mailing list or even a Slack-forwarding address) and the alert types they should subscribe to.
In the above notification method, I selected Data Movement Stalled here, with my current prefix of /sample_data/source-http-ingest.
Output:
For the alert type Data Movement Stalled, a default interval can be configured for a catalog prefix, while individual captures and materializations can override the interval from their Alerts tab. Estuary currently allows task-level Data Movement intervals between 1 hour and 7 days.
Beyond Estuary's built-in alerting capabilities, you can also configure your own alerts in platforms like Datadog or Prometheus which we will explore later in this tutorial (See: Method #3).
In your Dashboard, for each of your Sources, Collections, and Destinations you can see an Active Alerts Summary that shows if there is a Task Failure, or if it's Chronically Failing, or Idle.
Method #2: Using Estuary’s CLI tool - flowctl
For deeper troubleshooting, Estuary's flowctl CLI provides command-line access to task logs, runtime statistics, and alert configuration. It is useful when the dashboard indicates a problem and you need more detailed information about what a capture, derivation, or materialization is doing.
For more details on logs and task-level troubleshooting, see Estuary’s guide to troubleshooting a task with flowctl.
While the Estuary dashboard, provides a high-level view of your pipelines, flowctl is useful when you need more detailed operational information while developing, troubleshooting, or monitoring a flow. You can use it to inspect captures, collections, and materializations, view logs, examine statistics, and review alert subscriptions that help you understand both pipeline behavior and monitoring coverage.
The collected telemetry includes:
- Data volume: bytes and documents processed in and out
- Timing: task runtime, open time, and publish timestamps
- Shard activity: shard-level information such as keyBegin and rClockBegin
- Errors and warnings: error, warning, and failure information
- Task-level statistics: metrics for captures, derivations, and materializations
- Alert subscriptions: configured notification subscriptions and their scope across catalog prefixes or alert types
The flow is automatic: as a task executes, Estuary collects its statistics and makes them available through its logs and stats collections. This lets you inspect pipeline activity as it happens without instrumenting each task manually, while alert subscription information helps you verify how native monitoring and notifications are configured.
You can also configure alert thresholds and scope subscriptions to customize when and how you receive notifications by using flowctl, Estuary's CLI tool.
Make sure you install flowctl first:
plaintextbrew tap estuary/flowctl
brew install flowctlSetting up Alert Subscriptions
Here are the options for flowctl alert subscriptions:
- List - Lists all the alert subscriptions under the given prefix
- Subscribe - Subscribe to start receiving alerts, or alter an existing subscription to change which alert types it applies to
- Unsubscribe - Stop receiving alerts, either for a subset of alert types or everything
Subscriptions are scoped by catalog prefix, and a subscription can target a sub-prefix such as a single environment (for example acmeCo/prod/). To route environments differently, create a separate subscription for each prefix.
List
For example , lets start with listing current subscriptions:
plaintextflowctl alerts subscriptions listOutput:
Since we previously added /sample_data/source-http-ingest/ as a Prefix subscription, now it tracks various alert types and it shows here the details of creation of subscription. The default subscription is our user which does not have anything subscribed besides the free trial.
Next, we can subscribe to prefix using:
plaintextflowctl alerts subscriptions subscribe --prefix acmeCo/prod/ --email oncall@example.com --alert-type shard_failedSubscribing to an alert:
Subscribe an address to Task Failed alerts for RuheeShrestha/sample_data/ prefix.
Unsubscribing to an alert:
Here I remove a single alert type called ‘free trial’ from my email.
Configuring Alerts
Alerts that can be configured using flowctl alerts configs. Alert configurations are useful for reducing unnecessary notifications and adjusting monitoring sensitivity based on the expected behavior of a pipeline.
For example, a continuously updating CDC pipeline may need a relatively short data-stall threshold, while a webhook capture that receives data intermittently may need a longer threshold to avoid treating normal periods of inactivity as failures.
Alert conditions can be configured with: flowctl alerts configs update
1. Configure When Alerts Fire - dataMovementStalled
plaintextflowctl alerts configs update \ --prefix
RuheeShrestha/sample_data/ \ --set
dataMovementStalled.condition.stalledFor=2hIf webhook events are expected regularly, an alert can be configured when no new data arrives for two hours:
This provides a signal that the producer may have stopped sending events, the webhook integration may have been interrupted, or the upstream application may be unexpectedly inactive.
2. Tune Task-Failure Sensitivity - shardFailed
plaintextflowctl alerts configs update \ --prefix
RuheeShrestha/sample_data/ \ --set
shardFailed.condition.failures=5 \ --set
shardFailed.condition.per=2h3. Enable or Disable Alerts by Prefix - shardFailed.enabled
For example, Task Failed alerts can be enabled broadly for the tenant:
plaintextflowctl alerts configs update --prefix
RuheeShrestha/sample_data/ --set
shardFailed.enabled=falseMore specific prefixes override broader configurations field by field. This makes it possible to define organization-wide defaults while applying different monitoring policies to particular pipelines or environments.
Task Failure Alerts - taskChronicallyFailing
If the task keeps failing, the alert type may progress from "Task Failed" to "Task Chronically Failing." If the task remains in this chronically failing state and is unable to progress, the task may be disabled ("Task Auto-Disabled (Failing)") until you can address the root cause of the failure.
You can change the failure threshold per prefix by setting up configs for this:
Consider the task chronically failing after 24 hours
plaintextflowctl alerts configs update \
--prefix RuheeShrestha/sample_data/ \
--set taskChronicallyFailing.condition.failingFor=24hIdle Task Alerts - taskIdle
"Task Idle" alerts trigger when both of the following are true:
- The task has not processed any data for an extended time period
- The task has not been modified recently
We can set up the task to be idle when a task is expected to receive or process data periodically:.
plaintextflowctl alerts configs update \
--prefix RuheeShrestha/sample_data/ \
--set taskIdle.condition.idleFor=7dIf the task remains in this idle state, the task may be disabled ("Task Auto-Disabled (Idle)"). You may publish a new version of the task to keep it from being disabled or re-enable the task when you want to use it again.
Available threshold configurations include:
| Setting | Description | Default |
|---|---|---|
| dataMovementStalled.condition.stalledFor | Timeframe with no new data before alerting, e.g. 2h or 1d | |
| shardFailed.condition.failures | Failures in the last configured length of time | 3 |
| shardFailed.condition.per | Timeframe for task failures | 8h |
| taskChronicallyFailing.condition.failingFor | Timeframe a task has been repeatedly failing | 30d |
| taskIdle.condition.idleFor | Length of time a task has been idle | 30d |
Viewing Logs
For a complete view of logs, use flowctl in the command line and run the following on the task:
plaintextflowctl logs --task RuheeShrestha/sample_data/source-http-ingestTo view Logs from Last Hour:
plaintextflowctl logs --task RuheeShrestha/sample_data/source-http-ingest --since 1hTime options:
- -since 30m - Last 30 minutes
- -since 1h - Last hour
- -since 24h - Last 24 hours
- -since 7d - Last 7 days
To track real-time activity:
plaintextflowctl logs --task RuheeShrestha/sample_data/source-http-ingest --followTo monitor for errors:
plaintextflowctl logs --task RuheeShrestha/sample_data/source-http-ingest --level ERRORExpect Output: should be 0.
If the Capture hasn't ingested data in 30 minutes:
plaintext#Check recent logs
flowctl logs --task RuheeShrestha/sample_data/source-http-ingest --since 30m --level ERRORIf there are too many errors in the pipeline:
plaintext#See all errors from last hour
flowctl logs --task RuheeShrestha/sample_data/source-http-ingest --since 1h --level ERRORplaintext#Count error frequency
flowctl logs --task RuheeShrestha/sample_data/source-http-ingest --since 1h --level WARNViewing Stats
plaintext#View stats
flowctl raw stats --task RuheeShrestha/sample_data/source-http-ingest
#View Stats from Specific Time Range:
flowctl raw stats --task RuheeShrestha/sample_data/source-http-ingest --since 30mExample Output:
plaintext
{"_meta":{"uuid":"3a024919-ad49-11f1-8001-6551d8c87f43"},"capture":{"RuheeShrestha/sample_data/sample_data":{"lastPublishedAt":"2026-09-10T18:55:44.264207700+00:00","out":{"bytesTotal":827,"docsTotal":1},"right":{"bytesTotal":781,"docsTotal":1}}},"openSecondsTotal":0.0000974,"shard":{"build":"16029892200006cf","keyBegin":"00000000","kind":"capture","name":"RuheeShrestha/sample_data/source-http-ingest","rClockBegin":"00000000"},"ts":"2026-09-10T18:55:44.264110300+00:00","txnCount":1}The output for one transaction consists of the following stats:
Metadata & Identification
- _meta.uuid: Unique identifier for this specific shard execution
- ts: Timestamp of the event (Sept 10, 2026, 18:55:44 UTC)
- shard.name: RuheeShrestha/sample_data/source-http-ingest — your capture task's full path
Data Throughput
- out.docsTotal: 1 — Your producer sent 1 document in this cycle
- out.bytesTotal: 827 — That document was 827 bytes in size
- right.docsTotal: 1 — 1 document was successfully committed to Estuary (right = the destination/committed side)
- right.bytesTotal: 781 — Committed size is 781 bytes (slightly smaller due to Estuary's internal storage optimization)
Performance
- openSecondsTotal: 0.0000974 — The entire ingest operation took ~0.0001 seconds (extremely fast!)
- txnCount: 1 — This was 1 transaction
Method #3: Using the OpenMetrics API
Estuary's native tools provide immediate operational visibility and troubleshooting capabilities. When you need historical metric analysis, custom PromQL queries, dashboards, or more granular alerting, the same pipeline telemetry can be exposed through Estuary's OpenMetrics API and collected by systems such as Prometheus.
What Is the Estuary OpenMetrics API?
Estuary's OpenMetrics API exposes detailed metrics data on your captures, derivations, and materializations. This allows you to track your Estuary usage and pipeline activity in depth. Integrating the API with monitoring platforms like Prometheus or Datadog also allows you to implement alerts with greater specificity than Estuary currently offers natively.
You can use it to track changes in pipeline performance over time, build dashboards, and trigger alerts when a capture stops ingesting data, a materialization falls behind, or error rates increase.
Estuary exposes detailed metrics about pipeline activity, including:
- captured_*
- derived_*
- materialized_*
- logged_*
- _behind
Estuary currently tracks the following metrics in the OpenMetrics API. Metrics fall in one of two categories:
- Counters: cumulative values that increase over time, such as documents or bytes processed.
- Gauges: current point-in-time values that can increase or decrease, such as publication timestamps or bytes of backlog.
To authenticate the API, you will need an Estuary refresh token. You can generate one in the Admin panel of your dashboard.
- Create Refresh Token from Estuary Dashboard
- Go to https://dashboard.estuary.dev/admin/api
- Generate a refresh token
- Copy it (you'll paste it in Prometheus config)
- Your prefix: Whatever catalog path you use
Integrating OpenMetrics API to Prometheus
What Prometheus Does
Prometheus collects and stores Estuary’s OpenMetrics data as time-series metrics. It periodically scrapes the metrics endpoint, allowing you to query pipeline performance over time, track trends such as throughput and lag, and configure alerts when metrics exceed defined thresholds.
To use the OpenMetrics API with Prometheus, configure your prometheus.yml file to include Estuary's information. For example:
plaintextglobal:
scrape_interval: 1m
scrape_configs:
- job_name: estuary_metrics
scheme: https
bearer_token: REFRESH_TOKEN
metrics_path: /api/v1/metrics/PREFIX/
static_configs:
- targets: [api.estuary.dev]PREFIX: RuheeShrestha/sample_data/source-http-ingest
REFRESH_TOKEN: Use the one Estuary refresh token generated earlier.
This is a standard format for exposing metrics (Prometheus-compatible).
- Endpoint: https://api.estuary.dev/api/v1/metrics/{PREFIX}/
- Response: Plain text, one metric per line
- Scrape interval: 1 minute recommended
a. Setting up prometheus:
docker run --rm -it -p "9090:9090" -v $(pwd)/prometheus.yml:/etc/prometheus/prometheus.yml prom/prometheus:latest
Access Prometheus Dashboard on http://localhost:9090
b. Test the endpoint
plaintextcurl -H "Authorization: Bearer $TOKEN" \
"https://api.estuary.dev/api/v1/metrics/demo/" | \
grep "captured_in_docs_total"In Prometheus (http://localhost:9090/):
Checking Metrics from Captured Data
1. Total documents in a collection
captured_in_docs_total{collection="RuheeShrestha/sample_data/sample_data”}
Prometheus query result for captured_in_docs_total, showing the running total of documents captured into the sample_data collection. This is similar to the total documents being captured in Estuary’s UI/Dashboard Collections for each of your captures, or you could subscribe to the Data Movement Alerts.
The result shows 4,448 documents captured cumulatively for sample_data.
2. Total number of bytes used in a collection
captured_in_bytes_total{collection="RuheeShrestha/sample_data/sample_data”}
This counter tracks the cumulative size, in bytes, of all documents captured into the collection. This can easily be represented in the Collections for your sources in Estuary’s Dashboard as well.
3. Ingestion Activity
rate(captured_in_docs_total[5m]): captured_in_docs_total is a cumulative counter of pre-combine documents captured by the connector. Applying rate() converts the counter into the average number of documents captured per second over the previous five minutes.
The result is about 0.10 documents/sec, or roughly 6 documents/minute, for sample_data, while the other collections are inactive. This is a useful operational signal because it shows that records are actively arriving.
4. Throughput Over Time Window
increase(captured_in_docs_total[1h]): increase() reports the total growth of the counter over the specified window — here, how many documents were captured in the last hour.
The result is approximately 28.5 documents over the previous hour. Prometheus may return a fractional value because increase() extrapolates between sampled counter values, even though documents themselves are discrete.
5. Data Freshness
captured_last_published_at_time_seconds{...}: This gauge exposes the Unix timestamp of the most recently published document in the collection.
time() - captured_last_published_at_time_seconds{collection='RuheeShrestha/sample_data/sample_data'}:
Subtracting the last-published timestamp from the current time (time()) gives the age since the most recent captured publication.
The above query returns about 619 seconds, or 10.3 minutes since the last publication, since this result is already above the threshold of ~5 minutes it could be considered stale.
6. Tracking Errors
Prometheus queries also expose counters for task logs. To identify whether any errors occurred during the previous five minutes:
- Error spike: increase(logged_errors_total[5m]) > 0
This alert condition fires when any errors are logged in the trailing 5-minute window. Since a healthy pipeline should log at or near zero errors, any positive value here is worth investigating — check flowctl logs --level ERROR for the underlying cause.
- No-ingestion condition: rate(captured_in_docs_total[5m]) == 0
This alert condition fires when the ingestion rate drops to exactly zero over the trailing 5-minute window — a strong signal that the capture has stopped receiving data entirely, as opposed to merely slowing down.
This metric identifies logged errors; it should not be confused with Estuary's native Task Failed alerts, which monitor task-failure conditions directly.
These metrics provide detailed visibility into capture behavior. However, continuously watching dashboards or writing Prometheus rules for every operational condition is not always necessary. Estuary's native alerting system provides built-in notifications for common conditions such as task failures, stalled data movement, chronically failing tasks, and long-running idle tasks.
Checking Metrics from Materialization
Next we can check the materialized data from the collections into Google Sheets and check the metrics for it.
- Click on your collections and Click on Materialize.
- Add your Google Sheets link into the Spreadsheet URL
- Authenticate your Google Sheets
- Add your sheet name: webhook_data
- The data then gets loaded into the Google Sheets.
You can then track your Materialization metrics in Prometheus:
Materialization freshness
materialized_last_source_published_at_time_seconds: Reports the publication timestamp, in Unix seconds, of the most recent source-collection document that has been materialized to the destination. For this example, it indicates the publication time of the latest source document that has been processed by the Google Sheets materialization.
Prometheus query result for materialized_last_source_published_at_time_seconds, showing the timestamp of the most recent document written to the Google Sheets destination:
The result returns 1789092509..., corresponding to roughly 02:08:29 UTC, while Prometheus is evaluating around 02:21:41. That means the destination is roughly 13 minutes behind the latest source publication processed by the materialization.
Materialization Backlog
materialized_bytes_behind: is a gauge metric that measures how many bytes of source collection data a materialization has not yet processed for a given binding. It essentially represents the materialization’s current backlog.
Prometheus graph of materialized_bytes_behind, showing bytes read from the source that have not yet landed in the destination:
The result shows 12,867,917 bytes (~12.9 MB) of source-journal backlog still waiting to be read by the materialization
Other queries to test:
- materialized_in_docs_total — Cumulative pre-reduce documents read from the source collection by the materialization.
- materialized_out_docs_total — Cumulative post-reduce documents written to the destination.
- materialized_out_bytes_total — Cumulative post-reduce bytes stored in the destination.
- rate(materialized_out_docs_total[5m]) — Destination write throughput in documents per second over the last 5 minutes. Compare it with rate(captured_in_docs_total[5m]) to help distinguish source slowdowns from destination bottlenecks.
You can also view these metrics directly in the Materialization Overview section of the Estuary dashboard. The Overview provides a quick summary of materialization activity, including the total documents and total bytes read or written, so Prometheus is most useful when you want historical trends, custom rate calculations, or more detailed alerting.
Next Steps: Build Dashboards in Grafana
Once Prometheus is collecting metrics from Estuary’s OpenMetrics API, the next step is to visualize those metrics in Grafana. Grafana can connect directly to Prometheus and turn the collected time-series data into dashboards that provide a real-time view of pipeline health.
Note: Estuary’s OpenMetrics API primarily exposes operational pipeline metrics such as throughput, freshness, backlog, errors, and usage. Some of the dashboards below—particularly data-quality and infrastructure-efficiency views—require additional metrics to be collected from downstream systems or custom monitoring logic.
Dashboard 1 - Pipeline Health in Real Time
A pipeline health dashboard can provide a real-time view of whether data is moving through the pipeline as expected. The most useful panels focus on freshness, throughput, backlog, and errors, allowing you to identify both sudden failures and gradual performance degradation.
A real-time pipeline health dashboard should provide a compact view of the main operational signals across captures and materializations. Useful panels include:
- Freshness: time since the latest published record
- Capture throughput: documents ingested per second
- Materialization throughput: documents written downstream per second
- Backlog: bytes waiting to be processed by a derivation or materialization
- Errors and warnings: recent increases in logged errors, warnings, or failures
These panels are most useful when viewed together. For example, healthy capture throughput combined with increasing materialization backlog can indicate a downstream bottleneck, while a drop in both capture and materialization throughput may point to an upstream ingestion issue.
Dashboard 2 - Extending Observability with Data Quality Metrics
Pipeline health does not necessarily mean that the data itself is correct. A data-quality dashboard can extend operational observability with custom downstream metrics such as row counts, schema changes, and missing-value rates. These metrics are not automatically exposed through Estuary’s OpenMetrics API and would require additional instrumentation or validation at the collection or destination level.
This is particularly important when data is consumed by BI dashboards, analytics workflows, and machine learning models, where poor-quality or incomplete data can directly affect business decisions and model performance.
- Row counts: Compare expected versus actual row counts in Google Sheets
- Schema drift: Track column additions, removals, or data type changes
- Completeness: Monitor NULL or missing-value rates for individual columns
Dashboard 3 - Usage and Efficiency
Estuary also exposes usage and processing metrics that can help evaluate pipeline efficiency. For example, usage_seconds_total tracks billable connector usage time by task, while document and byte counters can be used to compare usage with pipeline throughput.
- Connector usage: Track usage_seconds_total over time
- Processing volume: Monitor documents and bytes processed
- Efficiency: Compare usage trends with throughput or data volume
Operational metrics help engineers maintain reliable pipelines, data-quality metrics help analysts and downstream consumers trust the data, and cost and efficiency metrics help teams scale their data infrastructure responsibly. This creates a foundation that supports a wide range of use cases, from BI and analytics to machine learning and operational decision-making.
Troubleshooting Checklist: Is It the Source, the Pipeline, or the Destination?
- Check source ingestion:
- Check the Estuary Capture Overview to confirm that documents and bytes are still being written. Monitor docsTotal and bytesTotal per collection per time window
- Check the capture’s Alerts tab for Data Movement Stalled or other active notifications.
- Check current ingestion in Prometheus: rate(captured_in_docs_total[5m])
- If the rate is 0 when data is expected, investigate the source producer, connectivity, or capture.
- Is it the Estuary pipeline?
- Check the Alerts tab in the Estuary dashboard for Task Failed, Task Chronically Failing, Data Movement Stalled, or auto-disable notifications.
- Check the Logs tab for warnings, failures, or connector errors.
- Use flowctl for more detailed troubleshooting:
flowctl logs --task <capture-task> --level ERROR: Use logs to identify the underlying cause of failures or stalled processing.
or inspect recent activity:
plaintext**flowctl logs --task <task-name> --since 1h**
- Check for recent errors in Prometheus:
increase(logged_errors_total[5m]): A value greater than 0 means errors were logged recently.
- Is it the destination?
- Check the Materialization Overview in Estuary for documents and bytes being read/written.
- Check the materialization’s Alerts tab for stalled data movement or task failures.
- In Prometheus:
- Check materialization backlog:
materialized_bytes_behind - A growing value means the materialization is falling behind.
- Check destination throughput: rate(materialized_out_docs_total[5m])
- Check destination freshness:
time() - materialized_last_source_published_at_time_seconds
When troubleshooting your pipeline in Estuary, compare capture and materialization behavior to identify where the issue is occurring. If both capture ingestion and materialization rates drop, investigate the source or capture. If capture throughput remains healthy while materialization backlog grows and destination throughput falls, the issue is more likely with the materialization or downstream destination. An increase in errors should be followed by a review of the Alerts tab and task logs, while a pipeline that appears to be running but exceeds its freshness SLA may indicate stalled or delayed data movement.
Start with the Estuary dashboard to review task status, active alerts, logs, document counts, and bytes processed. Use flowctl or external monitoring tools such as Prometheus when you need more detailed troubleshooting or historical analysis.
Best Practices for Observing Estuary Pipelines
Once Prometheus is collecting metrics from Estuary’s OpenMetrics API, focus on a small set of signals that reflect overall pipeline health.
- Data freshness: Track the time since the most recent record was published using captured/materialized_last_published_at_time_seconds. Set thresholds based on the expected source cadence and SLA, since a steadily increasing value may indicate stalled ingestion or delayed downstream processing.
- Throughput: Use rate(captured_in_docs_total[5m]) or rate(materialized_out_docs_total[5m]) to monitor documents processed per second. Compare current throughput with the pipeline’s normal baseline so sudden drops, spikes, or periods of inactivity are easier to identify.
- Backlog: Monitor materialized_bytes_behind or derived_bytes_behind to determine whether downstream processing is keeping pace with incoming data. A continuously increasing backlog can indicate that a derivation, materialization, or destination is falling behind.
- Errors: Track logged_errors_total, logged_warnings_total, and logged_failures_total to identify task-level issues. Correlate error spikes with changes in throughput, freshness, or backlog, then use Estuary logs or flowctl to investigate the underlying cause.
- Use native alerts first: Estuary’s built-in alerts cover common conditions such as stalled data movement, task failures, chronically failing tasks, and idle tasks. Use Prometheus or Datadog when you need custom thresholds, historical analysis, or more complex alert conditions.
- Correlate multiple signals: Avoid diagnosing pipeline health from a single metric. For example, healthy capture throughput combined with increasing materialization backlog usually points to a downstream issue rather than a source problem.
- Revisit thresholds as pipelines change: Update baselines, alert thresholds, and dashboards as source volumes, destinations, and business SLAs evolve. Monitoring configured during development may need to be adjusted for production workloads.
Wrapping Up
Observability is essential for real-time data pipelines because a task can appear to be running while data is delayed, stalled, or failing downstream. Tracking signals such as freshness, throughput, backlog, errors, and task status helps teams detect issues earlier, understand where a problem is occurring, and reduce the risk of stale or incomplete data reaching downstream users.
In this article, we built a layered observability workflow around an Estuary pipeline by using the dashboard for operational visibility and email notifications for specific alerting, flowctl to manage alert subscriptions and thresholds, and to inspect task logs and runtime statistics for deeper troubleshooting. For more granular monitoring, we exposed pipeline metrics through OpenMetrics API and scraped them into Prometheus where we captured various operational pipeline metrics for both source, and materialization to extend visibility from ingestion through downstream delivery.
Together, these tools support a source-to-destination troubleshooting workflow, helping teams isolate whether issues originate at ingestion, within the pipeline, or at the downstream destination while maintaining confidence in pipeline freshness and reliability.

About the author
Ruhee has a background in Computer Science and Economics and has worked as a Data Engineer for SaaS providing tech startups, where she has automated ETL processes using cutting-edge technologies and migrated data infrastructures to the cloud with AWS/Azure services. She is currently pursuing a Master’s in Business Analytics with a focus on Operations and AI at Worcester Polytechnic Institute.
























