不丢失#
- 生产
- 存储:持久化存储;Kafka实际上还通过副本机制,让持久化数据更可靠,即将每个主题划分为多个分区,并为每个分区创建多个副本。这确保了即使部分Broker发生故障,消息仍然可以从其他副本中恢复,并被消费者消费。
- 消费:另一种方式是手动提交,也就是由消费者业务代码自己控制提交时机,主动调用函数来提交。简单来说,为了保证消费最终被处理过,只有在消费端处理成功之后,才提交偏移到Broker,否则不进行偏移提交,这样下次拉取还能拉取到这条消息。
不重复#
生产:幂等性生产,对于每个 Topic Partition,Kafka生产者为每条消息分配一个递增的序列号。Kafka 该序列号是递增的,表示消息的顺序,Broker 会跟踪每个 Topic Partition 的最后一个已提交的序列号。
消费:
- Redis
- 唯一标识符:为每个消息分配一个全局唯一的标识符(如UUID)。
- Redis Set:将已消费的消息ID存储在Redis的Set数据结构中。每次消费消息前,检查该消息ID是否已存在于Set中。
- 原子操作:使用Redis的原子操作(如SISMEMBER和SADD)来检查和添加消息ID,确保操作的原子性。
- 如果存储用MySQL,如何进行幂等处理
MySQL常见实现幂等性的思路有如下三种:
- 唯一约束:在MySQL表中为消息ID创建一个唯一约束。
- INSERT IGNORE:使用 INSERT IGNORE 或 ON DUPLICATE KEY UPDATE 语句来尝试插入消息记录。如果消息ID已存在,则忽略或更新该记录。
- 事务:在需要的情况下,使用事务来确保多个操作的原子性。
- Redis
有序#
- 一种最简单的做法,就是根据业务确定分区,即每类业务自己一个分区,这样就可以实现业务消息有序。
- 业务内分区:我们前面说了,可以根据业务分区,而如果单个业务压力过大,我们就要考虑业务内再次切分。
- 子业务
- 客户
