文章目录 客户端发起请求 Broker处理请求 简单校验 获取分区号和元信息 构建返回数据 createResponse 问题 客户端发起请求 我们在分
文章目录
- 客户端发起请求
- Broker处理请求
- 简单校验
- 获取分区号和元信息
- 构建返回数据 createResponse
- 问题
客户端发起请求
我们在分析消费者的时候, 有看到调用FindCoordinatorRequest的请求
private RequestFuture<Void> sendFindCoordinatorRequest(Node node) {// initiate the group metadata request
log.debug("Sending FindCoordinator request to broker {}", node);
FindCoordinatorRequest.Builder requestBuilder =
new FindCoordinatorRequest.Builder(
new FindCoordinatorRequestData()
.setKeyType(CoordinatorType.GROUP.id())
.setKey(this.rebalanceConfig.groupId));
return client.send(node, requestBuilder)
.compose(new FindCoordinatorResponseHandler());
}
Broker处理请求
def handleFindCoordinatorRequest(request: RequestChannel.Request): Unit = {val findCoordinatorRequest = request.body[FindCoordinatorRequest]
// 根据协调器类型判断是否授权过
if (findCoordinatorRequest.data.keyType == CoordinatorType.GROUP.id &&
!authorize(request.context, DESCRIBE, GROUP, findCoordinatorRequest.data.key))
sendErrorResponseMaybeThrottle(request, Errors.GROUP_AUTHORIZATION_FAILED.exception)
else if (findCoordinatorRequest.data.keyType == CoordinatorType.TRANSACTION.id &&
!authorize(request.context, DESCRIBE, TRANSACTIONAL_ID, findCoordinatorRequest.data.key))
sendErrorResponseMaybeThrottle(request, Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED.exception)
else {
// get metadata (and create the topic if necessary)
val (partition, topicMetadata) = CoordinatorType.forId(findCoordinatorRequest.data.keyType) match {
case CoordinatorType.GROUP =>
val partition = groupCoordinator.partitionFor(findCoordinatorRequest.data.key)
val metadata = getOrCreateInternalTopic(GROUP_METADATA_TOPIC_NAME, request.context.listenerName)
(partition, metadata)
case CoordinatorType.TRANSACTION =>
val partition = txnCoordinator.partitionFor(findCoordinatorRequest.data.key)
val metadata = getOrCreateInternalTopic(TRANSACTION_STATE_TOPIC_NAME, request.context.listenerName)
(partition, metadata)
case _ =>
throw new InvalidRequestException("Unknown coordinator type in FindCoordinator request")
}
def createResponse(requestThrottleMs: Int): AbstractResponse = {
def createFindCoordinatorResponse(error: Errors, node: Node): FindCoordinatorResponse = {
new FindCoordinatorResponse(
new FindCoordinatorResponseData()
.setErrorCode(error.code)
.setErrorMessage(error.message)
.setNodeId(node.id)
.setHost(node.host)
.setPort(node.port)
.setThrottleTimeMs(requestThrottleMs))
}
val responseBody = if (topicMetadata.errorCode != Errors.NONE.code) {
createFindCoordinatorResponse(Errors.COORDINATOR_NOT_AVAILABLE, Node.noNode)
} else {
val coordinatorEndpoint = topicMetadata.partitions.asScala
.find(_.partitionIndex == partition)
.filter(_.leaderId != MetadataResponse.NO_LEADER_ID)
.flatMap(metadata => metadataCache.getAliveBroker(metadata.leaderId))
.flatMap(_.getNode(request.context.listenerName))
.filterNot(_.isEmpty)
coordinatorEndpoint match {
case Some(endpoint) =>
createFindCoordinatorResponse(Errors.NONE, endpoint)
case _ =>
createFindCoordinatorResponse(Errors.COORDINATOR_NOT_AVAILABLE, Node.noNode)
}
}
trace("Sending FindCoordinator response %s for correlation id %d to client %s."
.format(responseBody, request.header.correlationId, request.header.clientId))
responseBody
}
sendResponseMaybeThrottle(request, createResponse)
}
}
简单校验
根据协调器类型判断是否有被授权。协调器类型有 GROUP((byte) 0), TRANSACTION((byte) 1)两种
获取分区号和元信息
这里的接口分两种情况,一个是协调列席为GROUP 一个是 TRANSACTION
他们的处理逻辑都是一样的,只是处理的Topic不一样
GROUP 对应的Topic是 __consumer_offsets
TRANSACTION 对应的Topic是__transaction_state
这里我们主要分析一下 GROUP的情况
注意:创建这个Topic的的几个特殊属性:
属性
值
描述
cleanup.policy
compact
日志清理策略为 :紧缩
segment.bytes
10010241024
一个日志段的大小
compression.type
producer
压缩类型 为跟生产者保持一致
构建返回数据 createResponse
这里才是真正的找到协调器的主要逻辑, 这里的判断逻辑是
上面我们获取到的分区号是partition, 我们同样获取到了__consumer_offsets的元信息Metadata。
那我们就可以获取到这个分区号, 并且就能够找到该分区的LeaderId所属在哪个Broker上。
知道了哪个Broker, 那我们就能够获取到对应的EndPoint, 一个Broker可能同时有多个EndPoint(配置了多个监听器),那么我们应该使用哪个EndPoint呢?
这个的判断逻辑与上面说过的一样,客户端发起请求时候的监听器是哪个,那么这里就应该用哪个监听器。
注意:如果找到的分区Leader不存在 那么这个协调器就不存在
然后会返回异常:
The coordinator is not available
问题