Kafka源码分析15-消息发送完了内存怎么处理?
Kafka源码分析14-如何处理响应消息 Producer已经成功把消息发送出去了,本文分析生产者成功发送消息之后内存怎么处理?
直接看Sender 线程的completeBatch() 方法:
private void completeBatch(RecordBatch batch, ProduceResponse.PartitionResponse response, long correlationId, long now) { //如果处理成功那就是成功了,但是如果服务端那儿处理失败了,是不是也要给我们发送回来异常的信息。 //error 这个里面存储的就是服务端发送回来的异常码 Errors error = response.error; //如果响应里面带有异常 并且 这个请求是可以重试的 if (error != Errors.NONE && canRetry(batch, error)) { // retry log.warn("Got error produce response with correlation id {} on topic-partition {}, retrying ({} attempts left). Error: {}", correlationId, batch.topicPartition, this.retries - batch.attempts - 1, error); //重新把发送失败等着批次 加入到队列里面。 this.accumulator.reenqueue(batch, now); this.sensors.recordRetries(batch.topicPartition.topic(), batch.recordCount); } else { //这儿过来的数据:带有异常,但是不可以重试(1:压根就不让重试2:重试次数超了) //其余的都走这个分支。 RuntimeException exception; //如果响应里面带有 没有权限的异常 if (error == Errors.TOPIC_AUTHORIZATION_FAILED) //自己封装一个异常信息(自定义了异常) exception = new TopicAuthorizationException(batch.topicPartition.topic()); else exception = error.exception(); // tell the user the result of their request //TODO 核心代码 把异常的信息也给带过去了 //我们刚刚看的就是这儿的代码 //里面调用了用户传进来的回调函数 //回调函数调用了以后 //说明我们的一个完整的消息的发送流程就结束了。 batch.done(response.baseOffset, response.logAppendTime, exception); //看起来这个代码就是要回收资源的。 this.accumulator.deallocate(batch); if (error != Errors.NONE) this.sensors.recordErrors(batch.topicPartition.topic(), batch.recordCount); } if (error.exception() instanceof InvalidMetadataException) { if (error.exception() instanceof UnknownTopicOrPartitionException) log.warn("Received unknown topic or partition error in produce request on partition {}. The " + "topic/partition may not exist or the user may not have Describe access to it", batch.topicPartition); metadata.requestUpdate(); } // Unmute the completed partition. if (guaranteeMessageOrder) this.accumulator.unmutePartition(batch.topicPartition); }复制代码
this.accumulator.deallocate(batch);
这个就是对内存释放。
作者:hsfxuebao
链接:https://juejin.cn/post/7019558089472868366