Schedule icon
Script icon
Docker icon
If icon
SlackIncomingWebhook icon
Log icon
Return icon

MongoDB Oplog Window and Replication Lag Sentinel

Audits MongoDB replica sets for shrinking oplog windows and secondary sync lag, delivering actionable alerts to DBAs on Slack to avoid full resyncs.

Categories
DataInfrastructure

In MongoDB replica sets, the primary node records all write operations into a capped collection named the operations log (local.oplog.rs). Secondary members continuously tail and replay this oplog to maintain replica synchronicity. The oplog window represents the total duration of historical changes stored in this collection before older operations are overwritten.

When heavy write batches or bulk data migrations occur, the oplog fills up rapidly, shrinking the oplog window from days to just a few hours. If a secondary node experiences network hiccups, hardware pauses, or slow disk writes, its replication lag can easily exceed the primary's available oplog window. Once a secondary falls off the oplog horizon, it enters a degraded RECOVERING or STALE state and cannot catch up without a time-consuming and I/O-intensive initial sync from a primary snapshot.

This blueprint provides an automated database reliability guard. Scheduled daily, it inspects MongoDB replica set metadata, calculates remaining oplog window hours, measures replication lag across all secondary nodes, and dispatches actionable Slack alerts before nodes desynchronize.

How it works

  1. Scheduled Daily Audit: The daily_oplog_audit trigger (io.kestra.plugin.core.trigger.Schedule) initiates the audit every morning at 06:00 UTC.
  2. MongoDB Oplog and Replication Inspection: The scan_oplog_and_replication task (io.kestra.plugin.scripts.python.Script) inspects the timestamp difference between the oldest and newest oplog entries, measures member lag from replSetGetStatus, and exports mongodb_replication_report.json.
  3. Threshold Evaluation: The check_replication_risk flowable task (io.kestra.plugin.core.flow.If) branches based on whether the oplog window is dangerously narrow or if any secondary lag exceeds the alert threshold.
  4. Slack Alert Dispatch: When replication risk is detected, notify_slack_dba (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) delivers an incident card detailing the oplog capacity, the offending secondary node, and the exact replSetResizeOplog command.
  5. Healthy Confirmation: When all members are in sync and the oplog buffer is generous, log_healthy_replication records compliant status.
  6. Manifest Export: The export_replication_manifest task records execution state for database capacity planning.

Architecture diagram

flowchart TD
    A[Schedule: Daily 06:00 UTC] --> B[scan_oplog_and_replication: MongoDB local.oplog.rs API]
    B --> C{Replication Risk Detected?}
    C -- Yes --> D[notify_slack_dba: Slack Webhook]
    C -- No --> E[log_healthy_replication: Log]
    D --> F[export_replication_manifest: Return JSON]
    E --> F

Use cases

  • Secondary Node Desynchronization Prevention: Avoid multi-hour initial resyncs that saturate primary disk I/O and degrade application read throughput.
  • Bulk Write Window Governance: Verify that scheduled high-volume ETL jobs or data migrations do not exhaust the oplog buffer during execution.
  • Multi-Region Disaster Recovery Validation: Monitor cross-region read replicas to ensure geographical WAN latency does not compromise data redundancy SLAs.

Inputs

Name Type Default Description
mongodb_connection_uri STRING mongodb://mongodb-primary.internal:27017 Connection string to the MongoDB primary replica set node.
min_oplog_window_hours_warning FLOAT 24.0 Minimum acceptable oplog time capacity in hours.
max_allowed_replication_lag_seconds INT 300 Allowed secondary node replication lag before alert.
slack_channel STRING #dba-alerts Slack channel destination for replication notifications.

Expected outputs

  • {{ outputs.scan_oplog_and_replication.vars.has_replication_risk }}: Boolean flag indicating if replication risk was detected.
  • {{ outputs.scan_oplog_and_replication.vars.oplog_window_hours }}: Elapsed duration of the current oplog in hours.
  • {{ outputs.scan_oplog_and_replication.vars.lagging_nodes_count }}: Number of secondary nodes exceeding lag threshold.
  • {{ outputs.scan_oplog_and_replication.vars.top_lagging_node }}: Node hostname with the highest sync lag.
  • {{ outputs.scan_oplog_and_replication.outputFiles['mongodb_replication_report.json'] }}: Complete JSON diagnostic report.

Prerequisites

  • MongoDB credentials with clusterMonitor or read permissions on the local database and admin command access.
  • Incoming Webhook URL configured for your database administrator Slack channel.

Secrets

  • MONGODB_PASSWORD: MongoDB administrator password.
  • SLACK_WEBHOOK_URL: Slack Incoming Webhook endpoint URL.

Quick start

  1. Configure MONGODB_PASSWORD and SLACK_WEBHOOK_URL in your Kestra namespace secrets.
  2. Import this flow YAML into your Kestra workspace.
  3. Specify your replica set primary hostname in mongodb_connection_uri.
  4. Click Execute in the UI to perform an initial replication health probe across your replica set.

Common pitfalls and troubleshooting

  • Connecting to Secondaries: The oplog window calculation requires reading the primary node's local.oplog.rs collection; ensure the connection string routes to the active primary rather than secondary members.
  • Dynamic Oplog Resizing: Starting with MongoDB 3.6, the oplog can be resized dynamically online via replSetResizeOplog without restarting mongod processes.

How to extend

  • Add automated online resizing of the oplog collection via db.adminCommand({replSetResizeOplog: 1, size: new_size}) when storage volume margins allow.
  • Inspect WiredTiger dirty cache percentages during replication lag spikes to identify storage subsystem bottlenecks.
  • Route critical alerts to PagerDuty when secondary node lag exceeds 1 hour.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.