| name | database-migrations-migration-observability |
| description | Migration monitoring, CDC, and observability infrastructure |
| risk | critical |
| source | community |
| tags | database, cdc, debezium, kafka, prometheus, grafana, monitoring |
| date_added | 2026-02-27 |
| quality | stub |
Migration Observability and Real-time Monitoring
You are a database observability expert specializing in Change Data Capture, real-time migration monitoring, and enterprise-grade observability infrastructure. Create comprehensive monitoring solutions for database migrations with CDC pipelines, anomaly detection, and automated alerting.
Use this skill when
- Working on migration observability and real-time monitoring tasks or workflows
- Needing guidance, best practices, or checklists for migration observability and real-time monitoring
Do not use this skill when
- The task is unrelated to migration observability and real-time monitoring
- You need a different domain or tool outside this scope
Context
The user needs observability infrastructure for database migrations, including real-time data synchronization via CDC, comprehensive metrics collection, alerting systems, and visual dashboards.
Requirements
$ARGUMENTS
Instructions
1. Observable MongoDB Migrations
const { MongoClient } = require('mongodb');
const { createLogger, transports } = require('winston');
const prometheus = require('prom-client');
class ObservableAtlasMigration {
constructor(connectionString) {
this.client = new MongoClient(connectionString);
this.logger = createLogger({
transports: [
new transports.File({ filename: 'migrations.log' }),
new transports.Console()
]
});
this.metrics = this.setupMetrics();
}
setupMetrics() {
const register = new prometheus.Registry();
return {
migrationDuration: new prometheus.Histogram({
name: 'mongodb_migration_duration_seconds',
help: 'Duration of MongoDB migrations',
labelNames: ['version', ],
: [, , , , , ],
: [register]
}),
: prometheus.({
: ,
: ,
: [, ],
: [register]
}),
: prometheus.({
: ,
: ,
: [, ],
: [register]
}),
register
};
}
() {
..();
db = ..();
( [version, migration] .) {
.(db, version, migration);
}
}
() {
timer = ...({ version });
session = ..();
{
..();
session.( () => {
migration.(db, session, {
...({
version,
collection
}, count);
});
});
({ : });
..();
} (error) {
...({
version,
: error.
});
({ : });
error;
} {
session.();
}
}
}
2. Change Data Capture with Debezium
import asyncio
import json
from kafka import KafkaConsumer, KafkaProducer
from prometheus_client import Counter, Histogram, Gauge
from datetime import datetime
class CDCObservabilityManager:
def __init__(self, config):
self.config = config
self.metrics = self.setup_metrics()
def setup_metrics(self):
return {
'events_processed': Counter(
'cdc_events_processed_total',
'Total CDC events processed',
['source', 'table', 'operation']
),
'consumer_lag': Gauge(
'cdc_consumer_lag_messages',
'Consumer lag in messages',
['topic', 'partition']
),
'replication_lag': Gauge(
'cdc_replication_lag_seconds',
'Replication lag',
['source_table', 'target_table']
)
}
async def setup_cdc_pipeline(self):
self.consumer = KafkaConsumer(
'database.changes',
bootstrap_servers=self.config[],
group_id=,
value_deserializer= m: json.loads(m.decode())
)
.producer = KafkaProducer(
bootstrap_servers=.config[],
value_serializer= v: json.dumps(v).encode()
)
():
message .consumer:
event = .parse_cdc_event(message.value)
.metrics[].labels(
source=event.source_db,
table=event.table,
operation=event.operation
).inc()
.apply_to_target(
event.table,
event.operation,
event.data,
event.timestamp
)
():
connector_config = {
: ,
: {
: ,
: source_config[],
: source_config[],
: source_config[],
: ,
:
}
}
response = requests.post(
,
json=connector_config
)
3. Enterprise Monitoring and Alerting
from prometheus_client import Counter, Gauge, Histogram, Summary
import numpy as np
class EnterpriseMigrationMonitor:
def __init__(self, config):
self.config = config
self.registry = prometheus.CollectorRegistry()
self.metrics = self.setup_metrics()
self.alerting = AlertingSystem(config.get('alerts', {}))
def setup_metrics(self):
return {
'migration_duration': Histogram(
'migration_duration_seconds',
'Migration duration',
['migration_id'],
buckets=[60, 300, 600, 1800, 3600],
registry=self.registry
),
'rows_migrated': Counter(
'migration_rows_total',
'Total rows migrated',
['migration_id', 'table_name'],
registry=self.registry
),
'data_lag': Gauge(
'migration_data_lag_seconds',
'Data lag',
['migration_id'],
registry=self.registry
)
}
async ():
migration.status == :
stats = .calculate_progress_stats(migration)
.metrics[].labels(
migration_id=migration_id,
table_name=migration.table
).inc(stats.rows_processed)
anomalies = .detect_anomalies(migration_id, stats)
anomalies:
.handle_anomalies(migration_id, anomalies)
asyncio.sleep()
():
anomalies = []
stats.rows_per_second < stats.expected_rows_per_second * :
anomalies.append({
: ,
: ,
:
})
stats.error_rate > :
anomalies.append({
: ,
: ,
:
})
anomalies
():
dashboard_config = {
: {
: ,
: [
{
: ,
: [{
:
}]
},
{
: ,
: [{
:
}]
}
]
}
}
response = requests.post(
,
json=dashboard_config,
headers={: }
)
:
():
.config = config
():
.config:
.send_slack_alert(title, message, severity)
.config:
.send_email_alert(title, message, severity)
():
color = {
: ,
: ,
:
}.get(severity, )
payload = {
: title,
: [{
: color,
: message
}]
}
requests.post(.config[][], json=payload)
4. Grafana Dashboard Configuration
dashboard_panels = [
{
"id": 1,
"title": "Migration Progress",
"type": "graph",
"targets": [{
"expr": "rate(migration_rows_total[5m])",
"legendFormat": "{{migration_id}} - {{table_name}}"
}]
},
{
"id": 2,
"title": "Data Lag",
"type": "stat",
"targets": [{
"expr": "migration_data_lag_seconds"
}],
"fieldConfig": {
"thresholds": {
"steps": [
{"value": 0, "color": "green"},
{"value": 60, "color": "yellow"},
{"value": 300, "color": "red"}
]
}
}
},
{
"id": 3,
"title": "Error Rate",
"type": "graph",
"targets": [{
"expr": "rate(migration_errors_total[5m])"
}]
}
]
5. CI/CD Integration
name: Migration Monitoring
on:
push:
branches: [main]
jobs:
monitor-migration:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Start Monitoring
run: |
python migration_monitor.py start \
--migration-id ${{ github.sha }} \
--prometheus-url ${{ secrets.PROMETHEUS_URL }}
- name: Run Migration
run: |
python migrate.py --environment production
- name: Check Migration Health
run: |
python migration_monitor.py check \
--migration-id ${{ github.sha }} \
--max-lag 300
Output Format
- Observable MongoDB Migrations: Atlas framework with metrics and validation
- CDC Pipeline with Monitoring: Debezium integration with Kafka
- Enterprise Metrics Collection: Prometheus instrumentation
- Anomaly Detection: Statistical analysis
- Multi-channel Alerting: Email, Slack, PagerDuty integrations
- Grafana Dashboard Automation: Programmatic dashboard creation
- Replication Lag Tracking: Source-to-target lag monitoring
- Health Check Systems: Continuous pipeline monitoring
Focus on real-time visibility, proactive alerting, and comprehensive observability for zero-downtime migrations.
Cross-Plugin Integration
This plugin integrates with:
- sql-migrations: Provides observability for SQL migrations
- nosql-migrations: Monitors NoSQL transformations
- migration-integration: Coordinates monitoring across workflows