kafka只对“已提交”的消息做有限度的持久化保证
- 已提交的消息:当kafka的若干个Broker成功接受到一条消息并写入日志文件后,它们会告诉生产者程序这条消息已成功提交。
- 有限度的持久化保存:多个Kafka Broker,至少有一个Broker存活
“消息丢失”案例
案例1:生产者程序丢失数据
Kafka的Producer是异步发送消息,如果你调用的是producer.send(mes),通常会立刻返回,但是你不能认为消息已经成功发送完成
解决方案:Producer永远要使用带有回调通知的发送API,使用producer.send(msg, callback),它能准确的告诉你消息是否提交成功了
案例2:消费者程序丢失数据
Consumer程序有个位移的概念,表示这个Consumer当前消费的Topic分区位置
Consumer程序开启多个线程异步处理消息,而Consumer程序自动的向前更新位移,如果某线程运行失败,负责的消息没有成功消费,而这时候位移更新了,那么这条消息对于Consumer而言是丢失了
解决方案:如果是多线程异步处理消费信息,Consumer程序不要开启自动提交位移,而是要应用程序提交位移
最佳实践
- 不要使用producer.send(msg),而要使用producer.send(msg, callback)
- 设置acks=all。表示所有的副本Broker都要接受消息,这条消息才算是已提交
- 设置retries为一个较大的值,对应Producer的自动重试。当出现网络抖动,发送消息失败。此时配置retries > 0能够自动重试消息发送
- 设置unclean.leader.election.enable=false。Broker端参数,不允许落后太多的Broker称为Leader
- 设置replication.factor >= 3。Broker端参数,最好多保存几份副本,防止消息丢失、冗余
- 设置min.insync.repliocas > 1。Broker端参数,控股之消息至少写入一个副本才能算”已提交“
- 确保 replication.factor > min.insync.replicas。如果两者相等,那么只要有一个副本挂了,整个分区都没法工作。推荐设置为replication.factor = min.insync.replicas + 1
- 确保消息消费完再提交。Consumer端有个参数 enable.auto.commit。设置为false,并采用手动提交位移的方式
小结

