HDFS-13665. [SBN read] Move RPC response serialization into Server.doResponse(). Contributed by Plamen Jeliazkov.
This commit is contained in:
parent
64b7cf59bd
commit
e27708c2da
@ -856,15 +856,15 @@ private class RpcCall extends Call {
|
||||
final Writable rpcRequest; // Serialized Rpc request from client
|
||||
ByteBuffer rpcResponse; // the response for this call
|
||||
|
||||
private RpcResponseHeaderProto bufferedHeader; // the response header
|
||||
private Writable bufferedRv; // the byte response
|
||||
private ResponseParams responseParams; // the response params
|
||||
private Writable rv; // the byte response
|
||||
|
||||
RpcCall(RpcCall call) {
|
||||
super(call);
|
||||
this.connection = call.connection;
|
||||
this.rpcRequest = call.rpcRequest;
|
||||
this.bufferedRv = call.bufferedRv;
|
||||
this.bufferedHeader = call.bufferedHeader;
|
||||
this.rv = call.rv;
|
||||
this.responseParams = call.responseParams;
|
||||
}
|
||||
|
||||
RpcCall(Connection connection, int id) {
|
||||
@ -885,12 +885,10 @@ private class RpcCall extends Call {
|
||||
this.rpcRequest = param;
|
||||
}
|
||||
|
||||
public void setBufferedHeader(RpcResponseHeaderProto header) {
|
||||
this.bufferedHeader = header;
|
||||
}
|
||||
|
||||
public void setBufferedRv(Writable rv) {
|
||||
this.bufferedRv = rv;
|
||||
void setResponseFields(Writable returnValue,
|
||||
ResponseParams responseParams) {
|
||||
this.rv = returnValue;
|
||||
this.responseParams = responseParams;
|
||||
}
|
||||
|
||||
@Override
|
||||
@ -924,9 +922,7 @@ public Void run() throws Exception {
|
||||
populateResponseParamsOnError(e, responseParams);
|
||||
}
|
||||
if (!isResponseDeferred()) {
|
||||
setupResponse(this, responseParams.returnStatus,
|
||||
responseParams.detailedErr,
|
||||
value, responseParams.errorClass, responseParams.error);
|
||||
setResponseFields(value, responseParams);
|
||||
sendResponse();
|
||||
} else {
|
||||
if (LOG.isDebugEnabled()) {
|
||||
@ -981,13 +977,11 @@ void doResponse(Throwable t) throws IOException {
|
||||
setupResponse(call,
|
||||
RpcStatusProto.FATAL, RpcErrorCodeProto.ERROR_RPC_SERVER,
|
||||
null, t.getClass().getName(), StringUtils.stringifyException(t));
|
||||
} else if (alignmentContext != null) {
|
||||
// rebuild response with state context in header
|
||||
RpcResponseHeaderProto.Builder responseHeader =
|
||||
call.bufferedHeader.toBuilder();
|
||||
alignmentContext.updateResponseState(responseHeader);
|
||||
RpcResponseHeaderProto builtHeader = responseHeader.build();
|
||||
setupResponse(call, builtHeader, call.bufferedRv);
|
||||
} else {
|
||||
setupResponse(call, call.responseParams.returnStatus,
|
||||
call.responseParams.detailedErr, call.rv,
|
||||
call.responseParams.errorClass,
|
||||
call.responseParams.error);
|
||||
}
|
||||
connection.sendResponse(call);
|
||||
}
|
||||
@ -3012,6 +3006,9 @@ private void setupResponse(
|
||||
headerBuilder.setRetryCount(call.retryCount);
|
||||
headerBuilder.setStatus(status);
|
||||
headerBuilder.setServerIpcVersionNum(CURRENT_VERSION);
|
||||
if (alignmentContext != null) {
|
||||
alignmentContext.updateResponseState(headerBuilder);
|
||||
}
|
||||
|
||||
if (status == RpcStatusProto.SUCCESS) {
|
||||
RpcResponseHeaderProto header = headerBuilder.build();
|
||||
@ -3038,12 +3035,6 @@ private void setupResponse(
|
||||
|
||||
private void setupResponse(RpcCall call,
|
||||
RpcResponseHeaderProto header, Writable rv) throws IOException {
|
||||
if (alignmentContext != null && call.bufferedHeader == null
|
||||
&& call.bufferedRv == null) {
|
||||
call.setBufferedHeader(header);
|
||||
call.setBufferedRv(rv);
|
||||
}
|
||||
|
||||
final byte[] response;
|
||||
if (rv == null || (rv instanceof RpcWritable.ProtobufWrapper)) {
|
||||
response = setupResponseForProtobuf(header, rv);
|
||||
|
Loading…
Reference in New Issue
Block a user