AWS DMS Continuous Replication for Hybrid Cloud Architectures

Building reliable continuous replication between on-premises databases and AWS using DMS — handling schema drift, network failures, and validation at scale.

#aws#dms#replication#migration
Cover image for the article: AWS DMS Continuous Replication for Hybrid Cloud Architectures

Hybrid cloud isn't a temporary state for most enterprises — it's a permanent architecture. When your on-premises Oracle database must stay authoritative for compliance reasons but your cloud-native applications need that data with sub-second freshness, continuous replication becomes critical infrastructure. Not a migration tool you run once — a permanent data pipeline that must never fall behind.

We built a continuous replication architecture using AWS DMS that keeps 14 source databases synchronized with their AWS targets, processing 2.8 billion change events daily with a 99.97% uptime SLA.

The Problem: Stale Data in Cloud Applications

Our hybrid architecture had cloud-native microservices reading from AWS databases that were refreshed via nightly ETL jobs. This 24-hour data staleness created cascading problems:

IssueBusiness ImpactFrequency
Customer sees outdated account balanceSupport tickets, trust erosion200+ daily
Inventory count mismatchOverselling, fulfillment failures50+ daily
Compliance report discrepanciesAudit findings, regulatory riskMonthly
Pricing updates delayedRevenue leakageEvery price change
Employee data out of syncAccess control gapsEvery org change

The business demanded real-time data freshness. The infrastructure team demanded reliability. We needed both.

Architecture: Multi-Source Continuous Replication

DMS Continuous Replication Architecture

The architecture connects 14 source databases through a fleet of DMS replication instances to their respective AWS targets. Each source-target pair runs as an independent task with its own CDC stream, allowing independent failure recovery.

Replication Topology

SourceTargetRecords/DayAvg LagInstance Type
Oracle (Orders)Aurora PostgreSQL450M1.2sdms.r6i.2xlarge
Oracle (Inventory)Aurora PostgreSQL280M0.8sdms.r6i.xlarge
Oracle (Customers)DynamoDB120M2.1sdms.r6i.xlarge
SQL Server (Finance)Aurora PostgreSQL85M1.5sdms.r6i.large
SQL Server (HR)Aurora PostgreSQL12M0.5sdms.r6i.large
PostgreSQL (Analytics)Redshift1.8B45sdms.r6i.4xlarge

DMS Task Configuration

// DMS task creation with optimized CDC settings
import { DatabaseMigrationService } from '@aws-sdk/client-database-migration-service';

const dms = new DatabaseMigrationService({ region: 'us-east-1' });

async function createReplicationTask(config: ReplicationConfig) {
  const taskSettings = {
    TargetMetadata: {
      TargetSchema: config.targetSchema,
      SupportLobs: true,
      LobChunkSize: 64,
      LimitedSizeLobMode: true,
      LobMaxSize: 32768,
      ParallelLoadThreads: config.parallelThreads || 8,
      ParallelLoadBufferSize: 500,
      BatchApplyEnabled: true,
      FullLoadMaxRowsToCompare: 10000,
    },
    Logging: {
      EnableLogging: true,
      LogComponents: [
        { Id: 'TRANSFORMATION', Severity: 'LOGGER_SEVERITY_INFO' },
        { Id: 'SOURCE_UNLOAD', Severity: 'LOGGER_SEVERITY_DEFAULT' },
        { Id: 'TARGET_LOAD', Severity: 'LOGGER_SEVERITY_DEFAULT' },
        { Id: 'SOURCE_CAPTURE', Severity: 'LOGGER_SEVERITY_INFO' },
        { Id: 'TARGET_APPLY', Severity: 'LOGGER_SEVERITY_INFO' },
      ],
    },
    ChangeProcessingTuning: {
      BatchApplyPreserveTransaction: true,
      BatchApplyTimeoutMin: 1,
      BatchApplyTimeoutMax: 30,
      BatchSplitSize: 0,
      MemoryLimitTotal: 1024,
      MemoryKeepTime: 60,
      StatementCacheSize: 50,
    },
    ErrorBehavior: {
      DataErrorPolicy: 'LOG_ERROR',
      EventErrorPolicy: 'IGNORE',
      TableErrorPolicy: 'SUSPEND_TABLE',
      RecoverableErrorStopRetryAfterThrottlingMax: true,
      RecoverableErrorThrottlingMax: 1800,
      RecoverableErrorInterval: 5,
      ApplyErrorDeletePolicy: 'IGNORE_RECORD',
      ApplyErrorInsertPolicy: 'LOG_ERROR',
      ApplyErrorUpdatePolicy: 'LOG_ERROR',
    },
    StreamBuffer: {
      StreamBufferCount: 3,
      StreamBufferSizeInMB: 8,
      CtrlStreamBufferSizeInMB: 5,
    },
  };

  return dms.createReplicationTask({
    ReplicationTaskIdentifier: `cdc-${config.sourceName}-to-${config.targetName}`,
    SourceEndpointArn: config.sourceArn,
    TargetEndpointArn: config.targetArn,
    ReplicationInstanceArn: config.instanceArn,
    MigrationType: 'cdc',
    CdcStartPosition: config.startPosition,
    TableMappings: JSON.stringify(config.tableMappings),
    ReplicationTaskSettings: JSON.stringify(taskSettings),
  });
}

Handling Schema Drift

On-premises databases evolve independently of cloud targets. Schema changes on the source that aren't handled gracefully cause replication failures. We built an automated schema drift detection and resolution pipeline:

// Schema drift detector and resolver
interface SchemaDrift {
  type: 'COLUMN_ADDED' | 'COLUMN_REMOVED' | 'COLUMN_TYPE_CHANGED' | 'TABLE_ADDED';
  source: { table: string; column?: string; dataType?: string };
  target: { table: string; column?: string; dataType?: string };
  autoResolvable: boolean;
  resolution?: string;
}

async function detectSchemaDrift(
  sourcePool: Pool,
  targetPool: Pool,
  mappings: TableMapping[]
): Promise<SchemaDrift[]> {
  const drifts: SchemaDrift[] = [];

  for (const mapping of mappings) {
    const sourceColumns = await getTableColumns(sourcePool, mapping.sourceTable);
    const targetColumns = await getTableColumns(targetPool, mapping.targetTable);

    // Detect added columns
    for (const col of sourceColumns) {
      if (!targetColumns.find((tc) => tc.name === mapColumnName(col.name, mapping))) {
        drifts.push({
          type: 'COLUMN_ADDED',
          source: { table: mapping.sourceTable, column: col.name, dataType: col.dataType },
          target: { table: mapping.targetTable },
          autoResolvable: isNullableOrHasDefault(col),
          resolution: `ALTER TABLE ${mapping.targetTable} ADD COLUMN ${mapColumnName(col.name, mapping)} ${mapDataType(col.dataType)};`,
        });
      }
    }

    // Detect type changes
    for (const col of sourceColumns) {
      const targetCol = targetColumns.find((tc) => tc.name === mapColumnName(col.name, mapping));
      if (targetCol && !isCompatibleType(col.dataType, targetCol.dataType)) {
        drifts.push({
          type: 'COLUMN_TYPE_CHANGED',
          source: { table: mapping.sourceTable, column: col.name, dataType: col.dataType },
          target: { table: mapping.targetTable, column: targetCol.name, dataType: targetCol.dataType },
          autoResolvable: isWideningChange(col.dataType, targetCol.dataType),
          resolution: generateAlterStatement(mapping.targetTable, targetCol.name, col.dataType),
        });
      }
    }
  }

  return drifts;
}

We run drift detection every 6 hours. Auto-resolvable drifts (column additions with defaults, type widenings) are applied automatically. Breaking changes trigger PagerDuty alerts for manual review.

Network Resilience: Handling VPN Failures

Direct Connect links fail. VPN tunnels flap. The replication pipeline must handle network interruptions gracefully without requiring full-load recovery:

Network Resilience Architecture

DMS maintains CDC positions in its replication slot/log reader. After a network interruption, it resumes from the last committed position. However, prolonged outages (>24 hours for Oracle) risk source log rotation, which would require a full reload.

Our mitigation strategy:

Outage DurationRecovery StrategyData Loss Risk
< 5 minutesAutomatic resume from CDC positionNone
5-60 minutesResume + validation of affected tablesNone
1-24 hoursResume + full table validation scanNone (logs retained)
> 24 hoursTargeted full-load of affected tablesNone (but slower)
#!/bin/bash
# monitor-replication-lag.sh — Alert on growing lag or stalled tasks

TASKS=$(aws dms describe-replication-tasks \
  --filters Name=replication-task-arn,Values=${TASK_ARNS} \
  --query 'ReplicationTasks[].{
    Id:ReplicationTaskIdentifier,
    Status:Status,
    CDCLatency:ReplicationTaskStats.CDCLatencySource,
    TablesLoading:ReplicationTaskStats.TablesLoading,
    TablesErrored:ReplicationTaskStats.TablesErrored
  }' --output json)

echo "$TASKS" | jq -r '.[] | select(.CDCLatency > 60) | .Id' | while read task; do
  echo "ALERT: Task $task has CDC latency > 60 seconds"
  aws cloudwatch put-metric-data \
    --namespace DMS/Custom \
    --metric-name CDCLagCritical \
    --value 1 \
    --dimensions TaskId=$task
done

echo "$TASKS" | jq -r '.[] | select(.TablesErrored > 0) | .Id' | while read task; do
  echo "CRITICAL: Task $task has errored tables"
  # Auto-resume errored tables
  aws dms reload-tables \
    --replication-task-arn "arn:aws:dms:us-east-1:${ACCOUNT}:task:${task}" \
    --tables-to-reload $(get_errored_tables $task)
done

Validation: Continuous Data Integrity Checks

Replication without validation is a trust exercise. We run continuous validation to catch data divergence before downstream applications surface incorrect results:

Validation TypeFrequencyScopeDetection Time
Row count comparisonEvery 5 minAll tables5 min
Checksum validationHourlyHigh-priority tables1 hour
Full row comparisonDailySampled 1%24 hours
Business rule validationReal-timeCritical fieldsSeconds

The DMS built-in validation (row count + value comparison) catches most issues. We supplement with business-rule validators for critical data:

-- Continuous validation query: order totals must match
WITH source_totals AS (
  SELECT order_id, SUM(line_amount) as total
  FROM source_orders_replica
  WHERE updated_at > NOW() - INTERVAL '1 hour'
  GROUP BY order_id
),
target_totals AS (
  SELECT order_id, SUM(line_amount) as total
  FROM target_orders
  WHERE updated_at > NOW() - INTERVAL '1 hour'
  GROUP BY order_id
)
SELECT s.order_id, s.total as source_total, t.total as target_total
FROM source_totals s
JOIN target_totals t ON s.order_id = t.order_id
WHERE ABS(s.total - t.total) > 0.01;

Cost Breakdown

ComponentMonthly CostNotes
DMS instances (6 total)$4,200Mixed sizes per workload
DMS storage (logs + buffer)$8002TB total across instances
Network transfer (DX)$1,2002.8B events/day
CloudWatch (monitoring)$180Custom metrics + dashboards
Lambda (drift detection)$45Runs every 6 hours
Total$6,425/mo

Compared to the nightly ETL alternative ($2,800/month but with 24h data lag), the $3,600 premium eliminates $180K/year in business impact from stale data.

Key Takeaways

  1. Continuous replication is infrastructure, not a project — treat it like you treat your database: monitored, maintained, and SLA-bound.
  2. Batch apply mode is essential for high-throughput sources — single-transaction apply mode can't keep up with 450M changes/day.
  3. Schema drift detection prevents silent failures — automate what you can, alert on what you can't.
  4. Validate continuously, not just at setup — data divergence is a when, not an if.
  5. Size instances for peak, not average — CDC lag during peak hours compounds and takes hours to recover from.

Hybrid cloud replication done well is invisible to application teams — they query AWS databases and get fresh data. Done poorly, it becomes the single most common root cause of data quality incidents. Invest in the monitoring and validation that makes the difference.

Comments

    No comments yet. Be the first to share your thoughts.