📖 数据密集型设计

分布式数据同步与复制策略

深入探讨数据同步机制与复制策略

一、分布式数据同步概述

分布式数据同步是保障多个节点数据一致性的关键技术,在数据密集型应用中,高效的数据同步机制是构建可靠分布式系统的基础。

二、数据同步模式

2.1 数据同步模式对比

同步模式 描述 一致性 性能 适用场景
同步复制 写入后立即同步 强一致性 金融交易
异步复制 写入后异步同步 最终一致性 日志、统计
半同步复制 写入后部分同步 介于强和最终之间 中等 电商订单
延迟同步 延迟一段时间后同步 弱一致性 数据分析

2.2 数据同步架构图

graph TD A[数据同步] --> B[同步复制] A --> C[异步复制] A --> D[半同步复制] A --> E[延迟同步] B --> B1[主节点] B1 --> B2[从节点1] B1 --> B3[从节点2] C --> C1[主节点] C1 --> C2[异步队列] C2 --> C3[从节点] D --> D1[主节点] D1 --> D2[部分同步] D1 --> D3[部分异步]

三、数据复制策略

3.1 复制拓扑结构

graph TD A[复制拓扑] --> B[主从复制] A --> C[链式复制] A --> D[环形复制] A --> E[星型复制] B --> B1[主] B1 --> B2[从1] B1 --> B3[从2] C --> C1[主] C1 --> C2[中间] C2 --> C3[从] D --> D1[节点1] D1 --> D2[节点2] D2 --> D3[节点3] D3 --> D1 E --> E1[中心] E1 --> E2[节点1] E1 --> E3[节点2] E1 --> E4[节点3]

3.2 复制策略实现

public class DataReplicationService
{
    public async Task ReplicateAsync(string sourceNode, List targetNodes, ReplicationMode mode)
    {
        var data = await _dataService.GetDataAsync(sourceNode);
        
        switch (mode)
        {
            case ReplicationMode.Synchronous:
                await ReplicateSynchronouslyAsync(targetNodes, data);
                break;
            case ReplicationMode.Asynchronous:
                await ReplicateAsynchronouslyAsync(targetNodes, data);
                break;
            case ReplicationMode.SemiSynchronous:
                await ReplicateSemiSynchronouslyAsync(targetNodes, data);
                break;
        }
    }
    
    private async Task ReplicateSynchronouslyAsync(List targetNodes, object data)
    {
        foreach (var node in targetNodes)
        {
            await _dataService.SetDataAsync(node, data);
        }
    }
    
    private async Task ReplicateAsynchronouslyAsync(List targetNodes, object data)
    {
        foreach (var node in targetNodes)
        {
            _backgroundQueue.QueueBackgroundWorkItem(async token =>
            {
                await _dataService.SetDataAsync(node, data);
            });
        }
    }
    
    private async Task ReplicateSemiSynchronouslyAsync(List targetNodes, object data)
    {
        var synchronousCount = (int)Math.Ceiling(targetNodes.Count * 0.5);
        
        for (int i = 0; i < synchronousCount; i++)
        {
            await _dataService.SetDataAsync(targetNodes[i], data);
        }
        
        for (int i = synchronousCount; i < targetNodes.Count; i++)
        {
            _backgroundQueue.QueueBackgroundWorkItem(async token =>
            {
                await _dataService.SetDataAsync(targetNodes[i], data);
            });
        }
    }
}

public enum ReplicationMode { Synchronous, Asynchronous, SemiSynchronous }

3.3 数据复制冲突解决

public class ReplicationConflictResolver
{
    public async Task ResolveConflictAsync(string key, List versions)
    {
        var resolvedVersion = Resolve(versions);
        
        await _dataService.SetDataAsync(key, resolvedVersion.Data);
        
        await _dataService.SetMetadataAsync(key, new DataMetadata
        {
            Version = resolvedVersion.Version,
            ResolvedAt = DateTime.UtcNow
        });
    }
    
    private DataVersion Resolve(List versions)
    {
        var latestVersion = versions.OrderByDescending(v => v.Timestamp).First();
        
        var hasConflict = versions.GroupBy(v => v.Data.GetHashCode()).Count() > 1;
        
        if (!hasConflict)
        {
            return latestVersion;
        }
        
        return ResolveWithStrategy(versions);
    }
    
    private DataVersion ResolveWithStrategy(List versions)
    {
        return _conflictResolutionStrategy.Resolve(versions);
    }
}

public class DataVersion
{
    public object Data { get; set; }
    public long Version { get; set; }
    public DateTime Timestamp { get; set; }
    public string SourceNode { get; set; }
}

四、数据同步机制

4.1 基于日志的同步

public class LogBasedSyncService
{
    public async Task SyncFromLogAsync(string sourceNode, string targetNode)
    {
        var lastSyncPosition = await _metadataService.GetLastSyncPositionAsync(targetNode);
        
        var logs = await _logService.GetLogsAsync(sourceNode, lastSyncPosition);
        
        foreach (var log in logs)
        {
            await ApplyLogAsync(targetNode, log);
        }
        
        await _metadataService.UpdateLastSyncPositionAsync(targetNode, logs.Last().Position);
    }
    
    private async Task ApplyLogAsync(string targetNode, DataLog log)
    {
        switch (log.OperationType)
        {
            case OperationType.Insert:
                await _dataService.InsertAsync(targetNode, log.Data);
                break;
            case OperationType.Update:
                await _dataService.UpdateAsync(targetNode, log.Data);
                break;
            case OperationType.Delete:
                await _dataService.DeleteAsync(targetNode, log.Key);
                break;
        }
    }
    
    public async Task StreamLogsAsync(string sourceNode, string targetNode)
    {
        _logService.Subscribe(sourceNode, async log =>
        {
            await ApplyLogAsync(targetNode, log);
        });
    }
}

public class DataLog
{
    public long Position { get; set; }
    public OperationType OperationType { get; set; }
    public string Key { get; set; }
    public object Data { get; set; }
    public DateTime Timestamp { get; set; }
}

public enum OperationType { Insert, Update, Delete }

4.2 基于CDC的数据同步

public class CdcSyncService
{
    public async Task StartSyncAsync(string sourceDatabase, string targetDatabase, List tables)
    {
        var connector = CreateCdcConnector(sourceDatabase);
        
        connector.OnDataChange += async (sender, args) =>
        {
            await ApplyChangeAsync(targetDatabase, args);
        };
        
        await connector.StartAsync(tables);
    }
    
    private ICdcConnector CreateCdcConnector(string database)
    {
        return database switch
        {
            "mysql" => new MySqlCdcConnector(),
            "postgresql" => new PostgreSqlCdcConnector(),
            "sqlserver" => new SqlServerCdcConnector(),
            _ => throw new ArgumentException("Unsupported database")
        };
    }
    
    private async Task ApplyChangeAsync(string targetDatabase, DataChangeEvent args)
    {
        await _dataService.ApplyChangeAsync(targetDatabase, args);
    }
    
    public async Task StopSyncAsync()
    {
        await _connector.StopAsync();
    }
}

public interface ICdcConnector
{
    event EventHandler OnDataChange;
    Task StartAsync(List tables);
    Task StopAsync();
}

public class DataChangeEvent
{
    public string TableName { get; set; }
    public OperationType OperationType { get; set; }
    public object Data { get; set; }
    public DateTime Timestamp { get; set; }
}

4.3 基于消息队列的数据同步

public class MessageQueueSyncService
{
    public async Task PublishDataChangeAsync(DataChangeEvent changeEvent)
    {
        var message = new SyncMessage
        {
            OperationType = changeEvent.OperationType,
            TableName = changeEvent.TableName,
            Data = changeEvent.Data,
            Timestamp = changeEvent.Timestamp
        };
        
        await _messageQueue.PublishAsync("data_sync", message);
    }
    
    public async Task SubscribeToDataChangesAsync(string targetDatabase)
    {
        await _messageQueue.SubscribeAsync("data_sync", async message =>
        {
            await ApplyChangeAsync(targetDatabase, message);
        });
    }
    
    private async Task ApplyChangeAsync(string targetDatabase, SyncMessage message)
    {
        switch (message.OperationType)
        {
            case OperationType.Insert:
                await _dataService.InsertAsync(targetDatabase, message.Data);
                break;
            case OperationType.Update:
                await _dataService.UpdateAsync(targetDatabase, message.Data);
                break;
            case OperationType.Delete:
                await _dataService.DeleteAsync(targetDatabase, message.Key);
                break;
        }
    }
}

public class SyncMessage
{
    public OperationType OperationType { get; set; }
    public string TableName { get; set; }
    public string Key { get; set; }
    public object Data { get; set; }
    public DateTime Timestamp { get; set; }
}

五、数据同步监控与管理

5.1 数据同步监控

public class DataSyncMonitor
{
    public async Task GetMetricsAsync()
    {
        var metrics = new SyncMetrics();
        
        var syncTasks = await _syncRepository.GetAllSyncTasksAsync();
        
        foreach (var task in syncTasks)
        {
            metrics.TotalTasks++;
            
            if (task.Status == SyncStatus.Running)
            {
                metrics.RunningTasks++;
            }
            else if (task.Status == SyncStatus.Failed)
            {
                metrics.FailedTasks++;
            }
            
            metrics.TotalLatency += task.Latency;
            metrics.TotalRecords += task.RecordsSynced;
        }
        
        metrics.AverageLatency = syncTasks.Count > 0 ? metrics.TotalLatency / syncTasks.Count : 0;
        
        return metrics;
    }
    
    public async Task MonitorAsync()
    {
        var metrics = await GetMetricsAsync();
        
        if (metrics.FailedTasks > 0)
        {
            await _alertService.SendAlert("数据同步任务失败", 
                $"失败任务数: {metrics.FailedTasks}");
        }
        
        if (metrics.AverageLatency > 1000)
        {
            await _alertService.SendAlert("数据同步延迟过高", 
                $"平均延迟: {metrics.AverageLatency}ms");
        }
    }
}

public class SyncMetrics
{
    public int TotalTasks { get; set; }
    public int RunningTasks { get; set; }
    public int FailedTasks { get; set; }
    public double TotalLatency { get; set; }
    public double AverageLatency { get; set; }
    public long TotalRecords { get; set; }
}

public enum SyncStatus { Running, Paused, Failed, Completed }

5.2 数据同步管理

public class DataSyncManagementService
{
    public async Task CreateSyncTaskAsync(SyncTaskConfig config)
    {
        var task = new SyncTask
        {
            Id = Guid.NewGuid().ToString(),
            Source = config.Source,
            Target = config.Target,
            Tables = config.Tables,
            Mode = config.Mode,
            Status = SyncStatus.Paused,
            CreatedAt = DateTime.UtcNow
        };
        
        await _syncRepository.CreateAsync(task);
        
        return task;
    }
    
    public async Task StartSyncTaskAsync(string taskId)
    {
        var task = await _syncRepository.GetAsync(taskId);
        
        task.Status = SyncStatus.Running;
        
        await _syncRepository.UpdateAsync(task);
        
        await _syncService.StartSyncAsync(task);
    }
    
    public async Task PauseSyncTaskAsync(string taskId)
    {
        var task = await _syncRepository.GetAsync(taskId);
        
        task.Status = SyncStatus.Paused;
        
        await _syncRepository.UpdateAsync(task);
        
        await _syncService.PauseSyncAsync(task);
    }
    
    public async Task ResumeSyncTaskAsync(string taskId)
    {
        await StartSyncTaskAsync(taskId);
    }
    
    public async Task DeleteSyncTaskAsync(string taskId)
    {
        await _syncService.StopSyncAsync(taskId);
        
        await _syncRepository.DeleteAsync(taskId);
    }
    
    public async Task GetSyncTaskAsync(string taskId)
    {
        return await _syncRepository.GetAsync(taskId);
    }
    
    public async Task> GetAllSyncTasksAsync()
    {
        return await _syncRepository.GetAllAsync();
    }
}

六、数据同步最佳实践

6.1 同步模式选择

场景 推荐模式 原因
金融交易 同步复制 强一致性要求
电商订单 半同步复制 兼顾一致性和性能
日志统计 异步复制 高性能要求
数据分析 延迟同步 允许一定延迟

6.2 同步优化策略

public class SyncOptimizationService
{
    public void ConfigureBatchSync(int batchSize)
    {
        _syncService.BatchSize = batchSize;
    }
    
    public void ConfigureCompression(bool enabled)
    {
        _syncService.CompressionEnabled = enabled;
    }
    
    public void ConfigureConcurrency(int concurrency)
    {
        _syncService.Concurrency = concurrency;
    }
    
    public void ConfigureRetryPolicy(int maxRetries, TimeSpan backoff)
    {
        _syncService.RetryPolicy = new RetryPolicy(maxRetries, backoff);
    }
    
    public async Task OptimizeSyncTaskAsync(string taskId)
    {
        var task = await _syncRepository.GetAsync(taskId);
        
        var metrics = await _syncMonitor.GetTaskMetricsAsync(taskId);
        
        if (metrics.Latency > 1000)
        {
            ConfigureBatchSync(1000);
            ConfigureConcurrency(4);
        }
    }
}

6.3 冲突解决最佳实践

  • 使用版本号或时间戳检测冲突
  • 根据业务需求选择冲突解决策略
  • 记录冲突日志便于追溯
  • 提供手动解决冲突的能力
  • 定期校验数据一致性

七、总结

分布式数据同步与复制是保障多个节点数据一致性的关键技术。通过合理选择同步模式、实现高效的数据同步机制、做好监控和管理,能够构建可靠的分布式数据同步系统。定期校验数据一致性、优化同步性能,能够保障系统的稳定运行。