一、分布式消息系统概述
分布式消息系统是数据密集型应用的关键组件,通过异步消息传递实现系统解耦和流量削峰。事件驱动架构基于消息系统构建,通过事件的发布和订阅实现业务逻辑的松散耦合。
二、消息系统对比
2.1 主流消息系统对比
| 特性 | Kafka | RabbitMQ | Pulsar | RocketMQ |
|---|---|---|---|---|
| 吞吐量 | 极高 | 中等 | 极高 | 高 |
| 延迟 | 低 | 低 | 极低 | 低 |
| 持久化 | 磁盘 | 磁盘/内存 | 分层存储 | 磁盘 |
| 多租户 | 弱 | 弱 | 强 | 弱 |
| 协议支持 | 自定义 | AMQP | 多协议 | 自定义 |
2.2 事件驱动架构
graph TD
A[事件生产者] --> B[消息系统]
B --> C[事件存储]
B --> D[事件路由]
D --> E[事件消费者1]
D --> F[事件消费者2]
D --> G[事件消费者3]
E --> E1[业务处理]
E1 --> E2[发布新事件]
E2 --> B
F --> F1[业务处理]
F1 --> F2[发布新事件]
F2 --> B
G --> G1[业务处理]
G1 --> G2[发布新事件]
G2 --> B
三、Kafka消息系统
3.1 Kafka架构
graph TD
A[Producer] --> B[Topic]
B --> C[Partition 0]
B --> D[Partition 1]
B --> E[Partition 2]
C --> C1[Broker 1]
C --> C2[Broker 2]
C --> C3[Broker 3]
D --> D1[Broker 1]
D --> D2[Broker 2]
D --> D3[Broker 3]
E --> E1[Broker 1]
E --> E2[Broker 2]
E --> E3[Broker 3]
F[Consumer Group] --> C1
F --> D2
F --> E3
G[Consumer Group 2] --> C2
G --> D1
G --> E1
3.2 Kafka生产者实现
public class KafkaEventProducer
{
private readonly IProducer _producer;
public KafkaEventProducer(KafkaProducerConfig config)
{
var producerConfig = new ProducerConfig
{
BootstrapServers = config.BootstrapServers,
Acks = Acks.All,
Retries = 3,
BatchSize = 16384,
LingerMs = 5,
CompressionType = CompressionType.Snappy
};
_producer = new ProducerBuilder(producerConfig).Build();
}
public async Task ProduceAsync(string topic, string message)
{
var result = await _producer.ProduceAsync(topic, new Message
{
Value = message
});
if (result.Status != PersistenceStatus.Persisted)
{
throw new KafkaException($"消息发送失败: {result.Status}");
}
}
public async Task ProduceAsync(string topic, string key, string message)
{
var result = await _producer.ProduceAsync(topic, new Message
{
Key = key,
Value = message
});
}
public void Produce(string topic, string message, Action> callback)
{
_producer.Produce(topic, new Message { Value = message }, callback);
}
public void Flush(TimeSpan timeout)
{
_producer.Flush(timeout);
}
public void Dispose()
{
_producer.Dispose();
}
}
public class KafkaProducerConfig
{
public string BootstrapServers { get; set; }
public int Retries { get; set; } = 3;
public int BatchSize { get; set; } = 16384;
public int LingerMs { get; set; } = 5;
}
3.3 Kafka消费者实现
public class KafkaEventConsumer
{
private readonly IConsumer _consumer;
private readonly IEventHandler _eventHandler;
public KafkaEventConsumer(KafkaConsumerConfig config, IEventHandler eventHandler)
{
var consumerConfig = new ConsumerConfig
{
BootstrapServers = config.BootstrapServers,
GroupId = config.GroupId,
AutoOffsetReset = AutoOffsetReset.Earliest,
EnableAutoCommit = false,
MaxPollIntervalMs = 300000
};
_consumer = new ConsumerBuilder(consumerConfig).Build();
_eventHandler = eventHandler;
}
public async Task StartConsumingAsync(string topic, CancellationToken cancellationToken)
{
_consumer.Subscribe(topic);
while (!cancellationToken.IsCancellationRequested)
{
try
{
var result = _consumer.Consume(cancellationToken);
await _eventHandler.HandleAsync(result.Message.Value);
_consumer.Commit(result);
}
catch (ConsumeException ex)
{
_logger.LogError(ex, "消费消息失败");
}
}
}
public async Task StartConsumingBatchAsync(string topic, int batchSize, CancellationToken cancellationToken)
{
_consumer.Subscribe(topic);
var batch = new List>();
while (!cancellationToken.IsCancellationRequested)
{
try
{
var result = _consumer.Consume(cancellationToken);
batch.Add(result);
if (batch.Count >= batchSize)
{
await _eventHandler.HandleBatchAsync(batch.Select(r => r.Message.Value).ToList());
_consumer.Commit(batch.Last());
batch.Clear();
}
}
catch (ConsumeException ex)
{
_logger.LogError(ex, "消费消息失败");
}
}
}
public void Dispose()
{
_consumer.Dispose();
}
}
public class KafkaConsumerConfig
{
public string BootstrapServers { get; set; }
public string GroupId { get; set; }
public bool EnableAutoCommit { get; set; } = false;
public int MaxPollIntervalMs { get; set; } = 300000;
}
四、事件驱动架构实现
4.1 事件总线
public class EventBus
{
private readonly Dictionary>> _handlers = new();
private readonly IMessageProducer _messageProducer;
private readonly IMessageConsumer _messageConsumer;
public EventBus(IMessageProducer messageProducer, IMessageConsumer messageConsumer)
{
_messageProducer = messageProducer;
_messageConsumer = messageConsumer;
}
public void Subscribe(Func handler) where TEvent : Event
{
var eventName = typeof(TEvent).Name;
if (!_handlers.ContainsKey(eventName))
{
_handlers[eventName] = new List>();
}
_handlers[eventName].Add(e => handler((TEvent)e));
}
public void Unsubscribe(Func handler) where TEvent : Event
{
var eventName = typeof(TEvent).Name;
if (_handlers.ContainsKey(eventName))
{
_handlers[eventName].RemoveAll(h => h.Target == handler.Target && h.Method == handler.Method);
}
}
public async Task PublishAsync(TEvent @event) where TEvent : Event
{
@event.Id = Guid.NewGuid().ToString();
@event.Timestamp = DateTime.UtcNow;
var eventName = typeof(TEvent).Name;
var message = JsonSerializer.Serialize(@event);
await _messageProducer.ProduceAsync(eventName, message);
}
public async Task StartConsumingAsync(CancellationToken cancellationToken)
{
foreach (var eventName in _handlers.Keys)
{
await _messageConsumer.SubscribeAsync(eventName, async message =>
{
var @event = JsonSerializer.Deserialize(message);
if (@event != null)
{
await HandleEventAsync(@event);
}
});
}
await _messageConsumer.StartConsumingAsync(cancellationToken);
}
private async Task HandleEventAsync(Event @event)
{
var eventName = @event.GetType().Name;
if (_handlers.TryGetValue(eventName, out var handlers))
{
foreach (var handler in handlers)
{
await handler(@event);
}
}
}
}
public class Event
{
public string Id { get; set; }
public string EventName { get; set; }
public DateTime Timestamp { get; set; }
public string CorrelationId { get; set; }
}
4.2 事件溯源
public class EventSourcingService
{
private readonly IEventStore _eventStore;
private readonly Dictionary>> _eventHandlers = new();
public EventSourcingService(IEventStore eventStore)
{
_eventStore = eventStore;
}
public async Task AppendEventAsync(string aggregateId, Event @event)
{
@event.AggregateId = aggregateId;
@event.Version = await GetNextVersionAsync(aggregateId);
await _eventStore.AppendEventAsync(@event);
await PublishEventAsync(@event);
}
public async Task LoadAggregateAsync(string aggregateId) where T : IAggregateRoot, new()
{
var events = await _eventStore.GetEventsByAggregateIdAsync(aggregateId);
var aggregate = new T();
foreach (var @event in events)
{
aggregate.ApplyEvent(@event);
}
return aggregate;
}
public async Task> GetEventsAsync(string aggregateId)
{
return await _eventStore.GetEventsByAggregateIdAsync(aggregateId);
}
public async Task ReplayEventsAsync(string aggregateId)
{
var events = await _eventStore.GetEventsByAggregateIdAsync(aggregateId);
foreach (var @event in events)
{
await PublishEventAsync(@event);
}
}
private async Task GetNextVersionAsync(string aggregateId)
{
var events = await _eventStore.GetEventsByAggregateIdAsync(aggregateId);
return events.Count > 0 ? events.Max(e => e.Version) + 1 : 1;
}
private async Task PublishEventAsync(Event @event)
{
if (_eventHandlers.TryGetValue(@event.EventName, out var handlers))
{
foreach (var handler in handlers)
{
handler(@event);
}
}
await _eventBus.PublishAsync(@event);
}
}
public interface IAggregateRoot
{
void ApplyEvent(Event @event);
}
public interface IEventStore
{
Task AppendEventAsync(Event @event);
Task> GetEventsByAggregateIdAsync(string aggregateId);
Task> GetEventsByEventNameAsync(string eventName);
}
五、消息可靠性保障
5.1 消息可靠性策略
graph TD
A[发送消息] --> B{确认机制}
B -->|同步确认| C[等待Broker确认]
C -->|成功| D[返回成功]
C -->|失败| E[重试发送]
B -->|异步确认| F[立即返回]
F --> G[回调处理结果]
H[消费消息] --> I{消费模式}
I -->|自动提交| J[消费后自动提交offset]
I -->|手动提交| K[处理完成后手动提交]
K -->|成功| L[提交offset]
K -->|失败| M[重试消费]
M -->|多次失败| N[死信队列]
5.2 消息确认与重试
public class ReliableMessageProducer
{
private readonly IMessageProducer _producer;
private readonly int _maxRetries = 3;
private readonly TimeSpan _retryDelay = TimeSpan.FromSeconds(1);
public async Task SendWithRetryAsync(string topic, string message)
{
for (int attempt = 1; attempt <= _maxRetries; attempt++)
{
try
{
await _producer.ProduceAsync(topic, message);
return;
}
catch (Exception ex)
{
if (attempt == _maxRetries)
{
await HandleFailedMessageAsync(topic, message, ex);
throw;
}
await Task.Delay(_retryDelay * attempt);
}
}
}
private async Task HandleFailedMessageAsync(string topic, string message, Exception ex)
{
var deadLetterMessage = new DeadLetterMessage
{
Topic = topic,
Message = message,
Error = ex.Message,
RetryCount = _maxRetries,
Timestamp = DateTime.UtcNow
};
await _deadLetterQueue.EnqueueAsync(deadLetterMessage);
}
}
public class DeadLetterQueue
{
private readonly IMessageProducer _producer;
public async Task EnqueueAsync(DeadLetterMessage message)
{
var json = JsonSerializer.Serialize(message);
await _producer.ProduceAsync("dead-letter-queue", json);
}
public async Task DequeueAsync()
{
var message = await _consumer.ConsumeAsync("dead-letter-queue");
return JsonSerializer.Deserialize(message);
}
public async Task ProcessDeadLettersAsync()
{
while (true)
{
var message = await DequeueAsync();
await RetryMessageAsync(message);
}
}
private async Task RetryMessageAsync(DeadLetterMessage message)
{
if (message.RetryCount < 5)
{
await _producer.ProduceAsync(message.Topic, message.Message);
}
else
{
await _alertService.SendAlert("死信消息",
$"消息 {message.Message} 重试超过5次,已丢弃");
}
}
}
public class DeadLetterMessage
{
public string Topic { get; set; }
public string Message { get; set; }
public string Error { get; set; }
public int RetryCount { get; set; }
public DateTime Timestamp { get; set; }
}
5.3 幂等性保障
public class IdempotentMessageHandler
{
private readonly IIdempotentRepository _idempotentRepository;
private readonly TimeSpan _expiryTime = TimeSpan.FromHours(24);
public async Task HandleAsync(string messageId, Func> handler)
{
var exists = await _idempotentRepository.ExistsAsync(messageId);
if (exists)
{
var result = await _idempotentRepository.GetResultAsync(messageId);
if (result != null)
{
return result;
}
throw new DuplicateMessageException("消息已处理,但结果不存在");
}
await _idempotentRepository.LockAsync(messageId);
try
{
exists = await _idempotentRepository.ExistsAsync(messageId);
if (exists)
{
return await _idempotentRepository.GetResultAsync(messageId);
}
var result = await handler();
await _idempotentRepository.StoreResultAsync(messageId, result, _expiryTime);
return result;
}
finally
{
await _idempotentRepository.ReleaseLockAsync(messageId);
}
}
public async Task HandleAsync(string messageId, Func handler)
{
await HandleAsync(messageId, async () =>
{
await handler();
return true;
});
}
}
public interface IIdempotentRepository
{
Task ExistsAsync(string messageId);
Task GetResultAsync(string messageId);
Task StoreResultAsync(string messageId, T result, TimeSpan expiryTime);
Task LockAsync(string messageId);
Task ReleaseLockAsync(string messageId);
}
六、消息系统最佳实践
6.1 消息系统配置策略
| 配置项 | 建议值 | 说明 |
|---|---|---|
| Acks | all | 等待所有副本确认 |
| Retries | 3 | 重试次数 |
| Replication Factor | 3 | 副本数 |
| Min Insync Replicas | 2 | 最小同步副本数 |
| Auto Commit | false | 手动提交 |
6.2 事件驱动架构最佳实践
- 定义清晰的事件模型
- 实现消息幂等性
- 处理消息顺序
- 设置合理的超时时间
- 监控消息队列状态
七、总结
分布式消息系统与事件驱动架构是数据密集型应用的关键技术。通过选择合适的消息系统(Kafka、RabbitMQ、Pulsar等),能够实现系统解耦和流量削峰。事件驱动架构基于消息系统构建,通过事件的发布和订阅实现业务逻辑的松散耦合。消息可靠性保障(确认机制、重试、幂等性)是生产环境中必须考虑的问题。