消息传输过程图解
# 消息传输三个阶段
1.生产者把消息发送到MQ中
2.MQ收到消息并保存
3.消费者冲MQ拿到消息并进行消费
生产者数据可靠性配置
Kafka支持了ack机制,保证生产者可以将消息可靠的发送到达broker。
同时在消费消息的时候也有ack机制,来控制提交当前consumer的消费进度,上报offset
什么时候有可能丢消息:
1.没写进去acks=0
2.读出来了,没干活
生产者发送消息时可以配置ack的值:ProducerConfig.ACKS_CONFIG
acks=0:生产者发送过来数据,就不管了,可靠性差数据会丢,效率最高
acks=1:生产者发送过来数据,只需要Leader确认即可返回,可靠性中等,效率中等
acks=-1(all):生产者发送过来数据,Leader和ISR(所有在线的副本)队列里面所有Follwer应答,可靠性高效率最低
在生产环境中:
acks=0 很少使用;
acks=1,一般用于传输普通日志,允许丢个别数据;
acks=-1(all),一般用于传输重要不能丢失的数据(例如:钱、订单、积分等),对可靠性要求比较高的场景。
参考代码
public class CustomProducerAck { public static void main(string[] args) throws InterruptedException { //1.创建Kafka生产者的配置对象 Properties properties = new Properties(); // 2.给Kafka配置对象添加配置信息:bootstrap.servers properties.put(ProducerConfig.B00TSTRAP_SERVERS_coNFIG, "192.168.1.170:9092"); // key,value序列化(必须),key.serializer, value.serializerproperties.put(Producerconfig.KEY_sERIALIzER_CLASS_CONFIG, StringSerializer.class.getName (); properties. put(ProducerGonfig. VALUE_SERIALIzER_CLASs_coNFIG, StringSerializer.class.getName ()); //设置acks 可以设置"o","i"或者"all" properties.put(ProducerConfig.ACKS_CONFIG, "all"); //重试次数retries,默认Integer.MAX_VALUE (RETRY_BACKOFF_MS_CONFIG:重试间隔毫秒) properties.put(ProducerConfig.RETRIEs_CONFIG, 3); //3.创建Kafka生产者对象 KafkaProducer<string, String> kafkaProducer = new KafkaProducer<string, String> (properties); // 4.调用send方法,发送消息 for (int i= 0; i< 5; i++) { kafkaProducer.send(new ProducerRecord<>("second","atguigu "+ i)); } //5.关闭资源 kafkaProducer.close(); } }以上是生产端消息可靠性传输
此外还需要:
- MQ:保证MQ的高可用,一个分区创建多个副本,并将多个副本分散存储到kafkacluster的每一个节点上
- 消费者:将自动位移提交改为手动位移提交,保证消息消费完成
默认enable_auto_commit=true 5秒自动提交,关闭自动提交enable_auto_commit=false