diff --git a/core/src/main/java/com/taobao/arthas/core/shell/term/impl/http/HttpRequestHandler.java b/core/src/main/java/com/taobao/arthas/core/shell/term/impl/http/HttpRequestHandler.java index 011cdd97f..87876a5a2 100644 --- a/core/src/main/java/com/taobao/arthas/core/shell/term/impl/http/HttpRequestHandler.java +++ b/core/src/main/java/com/taobao/arthas/core/shell/term/impl/http/HttpRequestHandler.java @@ -69,10 +69,12 @@ public class HttpRequestHandler extends SimpleChannelInboundHandler byteBufPool = new ArrayBlockingQueue(poolSize); + private ArrayBlockingQueue charsBufPool = new ArrayBlockingQueue(poolSize); + private ArrayBlockingQueue bytesPool = new ArrayBlockingQueue(poolSize); + public static HttpApiHandler getInstance() { if (instance == null) { synchronized (HttpApiHandler.class) { @@ -71,11 +81,17 @@ public class HttpApiHandler { commandManager = sessionManager.getCommandManager(); jobController = sessionManager.getJobController(); historyManager = HistoryManagerImpl.getInstance(); + + //init buf pool + JsonUtils.setSerializeWriterBufferThreshold(jsonBufferSize); + for (int i = 0; i < poolSize; i++) { + byteBufPool.offer(Unpooled.buffer(jsonBufferSize)); + charsBufPool.offer(new char[jsonBufferSize]); + bytesPool.offer(new byte[jsonBufferSize]); + } } public HttpResponse handle(FullHttpRequest request) throws Exception { - DefaultFullHttpResponse response = new DefaultFullHttpResponse(request.protocolVersion(), - HttpResponseStatus.OK); ApiResponse result; String requestBody = null; @@ -97,12 +113,79 @@ public class HttpApiHandler { if (result == null) { result = createResponse(ApiState.FAILED, "The request was not processed"); } - result.setRequestId(requestId); - String jsonResult = JSON.toJSONString(result); - response.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/json; charset=utf-8"); - response.content().writeBytes(jsonResult.getBytes("UTF-8")); - return response; + + + //http response content + ByteBuf content = null; + //fastjson buf + char[] charsBuf = null; + byte[] bytesBuf = null; + + try { + //apply response content buf first + content = byteBufPool.poll(2000, TimeUnit.MILLISECONDS); + if (content == null) { + throw new ApiException("get response content buf failure"); + } + + //apply fastjson buf from pool + charsBuf = charsBufPool.poll(); + bytesBuf = bytesPool.poll(); + if (charsBuf == null || bytesBuf == null) { + throw new ApiException("get json buf failure"); + } + JsonUtils.setSerializeWriterBufThreadLocal(charsBuf, bytesBuf); + + //create http response + DefaultFullHttpResponse response = new DefaultFullHttpResponse(request.protocolVersion(), + HttpResponseStatus.OK, content.retain()); + response.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/json; charset=utf-8"); + writeResult(response, result); + return response; + } catch (Exception e) { + //response is discarded + if (content != null) { + content.release(); + byteBufPool.offer(content); + } + throw e; + } finally { + //give back json buf to pool + JsonUtils.setSerializeWriterBufThreadLocal(null, null); + if (charsBuf != null) { + charsBufPool.offer(charsBuf); + } + if (bytesBuf != null) { + bytesPool.offer(bytesBuf); + } + } + } + + public void onCompleted(DefaultFullHttpResponse httpResponse) { + ByteBuf content = httpResponse.content(); + content.clear(); + if (content.capacity() == jsonBufferSize) { + if (!byteBufPool.offer(content)) { + content.release(); + } + } else { + //replace content ByteBuf + content.release(); + if (byteBufPool.remainingCapacity() > 0) { + byteBufPool.offer(Unpooled.buffer(jsonBufferSize)); + } + } + } + + private void writeResult(DefaultFullHttpResponse response, Object result) throws IOException { + ByteBufOutputStream out = new ByteBufOutputStream(response.content()); + try { + JSON.writeJSONString(out, result); + } catch (IOException e) { + logger.error("write json to response failed", e); + throw e; + } } private ApiRequest parseRequest(String requestBody) throws ApiException { @@ -110,9 +193,9 @@ public class HttpApiHandler { throw new ApiException("parse request failed: request body is empty"); } try { - //Object jsonRequest = JSON.parse(requestBody); - ObjectMapper objectMapper = new ObjectMapper(); - return objectMapper.readValue(requestBody, ApiRequest.class); + //ObjectMapper objectMapper = new ObjectMapper(); + //return objectMapper.readValue(requestBody, ApiRequest.class); + return JSON.parseObject(requestBody, ApiRequest.class); } catch (Exception e) { throw new ApiException("parse request failed: " + e.getMessage(), e); } @@ -223,6 +306,7 @@ public class HttpApiHandler { /** * Update session input status for all consumer + * * @param session * @param inputStatus */ @@ -282,7 +366,7 @@ public class HttpApiHandler { response.setSessionId(session.getSessionId()) .setBody(body); - if(!session.tryLock()){ + if (!session.tryLock()) { response.setState(ApiState.REFUSED) .setMessage("Another command is executing."); return response; @@ -293,7 +377,7 @@ public class HttpApiHandler { Job job = null; try { Job foregroundJob = session.getForegroundJob(); - if (foregroundJob != null){ + if (foregroundJob != null) { response.setState(ApiState.REFUSED) .setMessage("Another job is running."); logger.info("Another job is running, jobId: {}", foregroundJob.id()); @@ -302,8 +386,8 @@ public class HttpApiHandler { //distribute result message both to origin session channel and request channel by CompositeResultDistributor packingResultDistributor = new PackingResultDistributorImpl(session); - ResultDistributor resultDistributor = new CompositeResultDistributorImpl(packingResultDistributor, session.getResultDistributor()); - job = this.createJob(commandLine, session, resultDistributor); + //ResultDistributor resultDistributor = new CompositeResultDistributorImpl(packingResultDistributor, session.getResultDistributor()); + job = this.createJob(commandLine, session, packingResultDistributor); session.setForegroundJob(job); updateSessionInputStatus(session, InputStatus.ALLOW_INTERRUPT); @@ -364,7 +448,7 @@ public class HttpApiHandler { response.setSessionId(session.getSessionId()) .setBody(body); - if(!session.tryLock()){ + if (!session.tryLock()) { response.setState(ApiState.REFUSED) .setMessage("Another command is executing."); return response; @@ -373,7 +457,7 @@ public class HttpApiHandler { try { Job foregroundJob = session.getForegroundJob(); - if (foregroundJob != null){ + if (foregroundJob != null) { response.setState(ApiState.REFUSED) .setMessage("Another job is running."); logger.info("Another job is running, jobId: {}", foregroundJob.id()); @@ -403,7 +487,7 @@ public class HttpApiHandler { CommandRequestModel commandRequestModel = new CommandRequestModel(commandLine, response.getState(), response.getMessage()); session.getResultDistributor().appendResult(commandRequestModel); return response; - }finally { + } finally { if (session.getLock() == lock) { session.unLock(); } @@ -439,7 +523,7 @@ public class HttpApiHandler { } ResultConsumer consumer = session.getResultDistributor().getConsumer(consumerId); if (consumer == null) { - throw new ApiException("consumer not found: "+consumerId); + throw new ApiException("consumer not found: " + consumerId); } List results = consumer.pollResults(); @@ -510,7 +594,7 @@ public class HttpApiHandler { @Override public void onBackground(Job job) { - if (session.getForegroundJob() == job){ + if (session.getForegroundJob() == job) { session.setForegroundJob(null); updateSessionInputStatus(session, InputStatus.ALLOW_INPUT); } @@ -518,7 +602,7 @@ public class HttpApiHandler { @Override public void onTerminated(Job job) { - if (session.getForegroundJob() == job){ + if (session.getForegroundJob() == job) { session.setForegroundJob(null); updateSessionInputStatus(session, InputStatus.ALLOW_INPUT); } @@ -526,7 +610,7 @@ public class HttpApiHandler { @Override public void onSuspend(Job job) { - if (session.getForegroundJob() == job){ + if (session.getForegroundJob() == job) { session.setForegroundJob(null); updateSessionInputStatus(session, InputStatus.ALLOW_INPUT); } diff --git a/core/src/main/java/com/taobao/arthas/core/util/JsonUtils.java b/core/src/main/java/com/taobao/arthas/core/util/JsonUtils.java new file mode 100644 index 000000000..5ab6acb67 --- /dev/null +++ b/core/src/main/java/com/taobao/arthas/core/util/JsonUtils.java @@ -0,0 +1,93 @@ +package com.taobao.arthas.core.util; + +import com.alibaba.fastjson.serializer.SerializeWriter; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.lang.reflect.Field; + +/** + * @author gongdewei 2020/5/15 + */ +public class JsonUtils { + private static final Logger logger = LoggerFactory.getLogger(JsonUtils.class); + private static Field serializeWriterBufLocalField; + private static Field serializeWriterBytesBufLocal; + private static Field serializeWriterBufferThreshold; + + /** + * Set Fastjson SerializeWriter Buffer Threshold + * @param value + */ + public static void setSerializeWriterBufferThreshold(int value) { + Class clazz = SerializeWriter.class; + try { + if (serializeWriterBufferThreshold == null) { + serializeWriterBufferThreshold = clazz.getDeclaredField("BUFFER_THRESHOLD"); + } + serializeWriterBufferThreshold.setAccessible(true); + serializeWriterBufferThreshold.set(null, value); + } catch (Exception e) { + logger.error("update SerializeWriter.BUFFER_THRESHOLD value failed", e); + } + } + + /** + * Set Fastjson SerializeWriter ThreadLocal value + * @param bufSize + */ + public static void setSerializeWriterBufThreadLocal(int bufSize) { + Class clazz = SerializeWriter.class; + try { + //set threadLocal value + if (serializeWriterBufLocalField == null) { + serializeWriterBufLocalField = clazz.getDeclaredField("bufLocal"); + } + serializeWriterBufLocalField.setAccessible(true); + ThreadLocal bufLocal = (ThreadLocal) serializeWriterBufLocalField.get(null); + char[] charsLocal = bufLocal.get(); + if (charsLocal == null || charsLocal.length < bufSize) { + bufLocal.set(new char[bufSize]); + } + + if (serializeWriterBytesBufLocal == null) { + serializeWriterBytesBufLocal = clazz.getDeclaredField("bytesBufLocal"); + } + serializeWriterBytesBufLocal.setAccessible(true); + ThreadLocal bytesBufLocal = (ThreadLocal) serializeWriterBytesBufLocal.get(null); + byte[] bytesLocal = bytesBufLocal.get(); + if (bytesLocal == null || bytesLocal.length < bufSize) { + bytesBufLocal.set(new byte[bufSize]); + } + } catch (Exception e) { + logger.error("update SerializeWriter.BUFFER_THRESHOLD value failed", e); + } + } + + /** + * Set Fastjson SerializeWriter ThreadLocal value + */ + public static void setSerializeWriterBufThreadLocal(char[] charsBuf, byte[] bytesBuf) { + Class clazz = SerializeWriter.class; + try { + //set threadLocal value + if (serializeWriterBufLocalField == null) { + serializeWriterBufLocalField = clazz.getDeclaredField("bufLocal"); + } + serializeWriterBufLocalField.setAccessible(true); + ThreadLocal bufLocal = (ThreadLocal) serializeWriterBufLocalField.get(null); + bufLocal.set(charsBuf); + + if (serializeWriterBytesBufLocal == null) { + serializeWriterBytesBufLocal = clazz.getDeclaredField("bytesBufLocal"); + } + serializeWriterBytesBufLocal.setAccessible(true); + ThreadLocal bytesBufLocal = (ThreadLocal) serializeWriterBytesBufLocal.get(null); + bytesBufLocal.set(bytesBuf); + } catch (Exception e) { + logger.error("update SerializeWriter.BUFFER_THRESHOLD value failed", e); + } + } + + +}