New to Kestra?
Use blueprints to kickstart your first workflows.
Audits MongoDB replica sets for shrinking oplog windows and secondary sync lag, delivering actionable alerts to DBAs on Slack to avoid full resyncs.
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.
daily_oplog_audit trigger (io.kestra.plugin.core.trigger.Schedule) initiates the audit every morning at 06:00 UTC.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.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.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.log_healthy_replication records compliant status.export_replication_manifest task records execution state for database capacity planning.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
| 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. |
{{ 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.clusterMonitor or read permissions on the local database and admin command access.MONGODB_PASSWORD: MongoDB administrator password.SLACK_WEBHOOK_URL: Slack Incoming Webhook endpoint URL.MONGODB_PASSWORD and SLACK_WEBHOOK_URL in your Kestra namespace secrets.mongodb_connection_uri.local.oplog.rs collection; ensure the connection string routes to the active primary rather than secondary members.replSetResizeOplog without restarting mongod processes.db.adminCommand({replSetResizeOplog: 1, size: new_size}) when storage volume margins allow.