From 9c23c745c76157ff32ce4c00af0556ba45aaa652 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Fri, 15 Jul 2016 09:47:55 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E5=B0=86=E6=8E=A5=E5=8F=97?= =?UTF-8?q?=E5=88=B0=E7=9A=84=E6=95=B0=E6=8D=AE=E7=BC=93=E5=AD=98=E5=88=B0?= =?UTF-8?q?=E7=BC=93=E5=AD=98=E6=96=87=E4=BB=B6=E4=B8=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../context/CurrentThreadSpanStack.java | 2 +- .../invoke/monitor/BaseInvokeMonitor.java | 15 +- .../protocol/util/ContextGenerator.java | 2 +- .../ai/cloud/serialize/SerializeTest.java | 16 +- .../protocol/proto/TraceProtocol.java | 173 ++++++++++++- .../ai/cloud/skywalking/protocol/AckSpan.java | 14 +- .../skywalking/protocol/RequestSpan.java | 31 ++- .../protocol/SerializedFactory.java | 20 +- .../ai/cloud/skywalking/protocol/Span.java | 16 +- .../protocol/TransportPackager.java | 52 ++-- .../common/AbstractDataSerializable.java | 5 +- .../protocol/common/ISerializable.java | 4 +- .../exception/ConvertFailedException.java | 4 + .../src/main/proto/TraceProtocol.proto | 5 +- .../dubbo-plugin/pom.xml | 5 - .../skywalking/reciever/CollectionServer.java | 6 +- .../reciever/buffer/AppendEOFFlagThread.java | 44 ++++ .../reciever/buffer/DataBufferThread.java | 102 ++++++-- .../buffer/DataBufferThreadContainer.java | 39 ++- .../skywalking/reciever/conf/Config.java | 29 +++ .../handler/CollectionServerDataHandler.java | 34 ++- .../model/BufferDataPackagerGenerator.java | 39 +++ .../peresistent/BufferBalePlucker.java | 227 ++++++++++++++++++ .../peresistent/FileRegisterEntry.java | 78 ++++++ .../reciever/peresistent/MemoryRegister.java | 127 ++++++++++ .../reciever/peresistent/OffsetFile.java | 69 ++++++ .../peresistent/PersistenceThread.java | 142 +++++++++++ .../PersistenceThreadLauncher.java | 12 + .../RegisterPersistenceThread.java | 98 ++++++++ .../reciever/processor/AckSpanProcessor.java | 1 + .../reciever/processor/DefaultProcessor.java | 12 + .../reciever/processor/ProcessorFactory.java | 12 +- .../processor/RequestSpanProcessor.java | 4 +- .../src/main/resources/config.properties | 24 +- .../peresistent/BufferBalePluckerTest.java | 74 ++++++ 35 files changed, 1397 insertions(+), 140 deletions(-) create mode 100644 skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/exception/ConvertFailedException.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/FileRegisterEntry.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/MemoryRegister.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/OffsetFile.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/RegisterPersistenceThread.java create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/DefaultProcessor.java create mode 100644 skywalking-server/src/test/java/test/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePluckerTest.java diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java index 5e85daafb..251c803d0 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java @@ -86,7 +86,7 @@ public class CurrentThreadSpanStack { public void invalidatePresentSpans() { for (SpanNode spanNode : spans) { - spanNode.getData().setInvalidate(true); + spanNode.getData().setValidate(true); } } } diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/BaseInvokeMonitor.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/BaseInvokeMonitor.java index b5dacbb26..4c1e990eb 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/BaseInvokeMonitor.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/BaseInvokeMonitor.java @@ -31,12 +31,13 @@ public abstract class BaseInvokeMonitor { logger.debug("TraceId:" + spanData.getTraceId() + "\tviewpointId:" + id.getViewPoint() + "\tParentLevelId:" + spanData.getParentLevel() + "\tLevelId:" + spanData .getLevelId()); } - // 根据SpanData生成RequestSpan,并保存 - ContextBuffer.save(RequestSpan.RequestSpanBuilder. - newBuilder(spanData).callType(id.getCallType()) - .viewPoint(id.getViewPoint()) - .spanTypeDesc(id.getSpanTypeDesc()) - .build()); + + + if (!spanData.isValidate()) { + // 根据SpanData生成RequestSpan,并保存 + ContextBuffer.save(RequestSpan.RequestSpanBuilder. + newBuilder(spanData).callType(id.getCallType()).viewPoint(id.getViewPoint()).spanTypeDesc(id.getSpanTypeDesc()).build()); + } // 将新创建的Context存放到ThreadLocal栈中。 CurrentThreadSpanStack.push(spanData); @@ -51,7 +52,7 @@ public abstract class BaseInvokeMonitor { // 弹出上下文的栈顶中的元素 Span spanData = CurrentThreadSpanStack.pop(); - if (spanData == null || spanData.isInvalidate()) { + if (spanData == null || spanData.isValidate()) { return; } diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/protocol/util/ContextGenerator.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/protocol/util/ContextGenerator.java index b032ee6fe..54eadc97b 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/protocol/util/ContextGenerator.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/protocol/util/ContextGenerator.java @@ -63,7 +63,7 @@ public final class ContextGenerator { // if true, current span is invalidate and current span also belong to RPC span if (parentSpan.isRPCClientSpan()) { span.setSpanType(SpanType.RPC_CLIENT); - span.setIsInvalidate(true); + span.setValidate(false); } if (!StringUtil.isEmpty(parentSpan.getParentLevel())) { diff --git a/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java b/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java index ea3221941..9aecd1e6a 100644 --- a/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java +++ b/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java @@ -2,19 +2,19 @@ package test.ai.cloud.serialize; import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.protocol.AckSpan; -import com.ai.cloud.skywalking.protocol.RequestSpan; import com.ai.cloud.skywalking.protocol.Span; -import com.ai.cloud.skywalking.protocol.common.CallType; import com.ai.cloud.skywalking.protocol.common.SpanType; public class SerializeTest { public static void main(String[] args) throws InterruptedException { - Span spandata = new Span("1.0b.1461060884539.7d6d06e.22489.1271.103", "", 0, "test-application", "test"); - spandata.setSpanType(SpanType.LOCAL); - spandata.setStartDate(System.currentTimeMillis() - 1000 * 60); - AckSpan requestSpan = new AckSpan(spandata); - ContextBuffer.save(requestSpan); + while (true) { + Span spandata = new Span("1.0b.1461060884539.7d6d06e.22489.1271.103", "", 0, "test-application", "test"); + spandata.setSpanType(SpanType.LOCAL); + spandata.setStartDate(System.currentTimeMillis() - 1000 * 60); + AckSpan requestSpan = new AckSpan(spandata); + ContextBuffer.save(requestSpan); + Thread.sleep(2 * 1000); + } - Thread.sleep(60 * 1000 * 10); } } diff --git a/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java b/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java index 1261262b5..af9be17c9 100644 --- a/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java +++ b/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java @@ -1206,6 +1206,20 @@ public final class TraceProtocol { */ com.google.protobuf.ByteString getBussinessKeyBytes(); + + /** + * required string agentId = 12; + */ + boolean hasAgentId(); + /** + * required string agentId = 12; + */ + java.lang.String getAgentId(); + /** + * required string agentId = 12; + */ + com.google.protobuf.ByteString + getAgentIdBytes(); } /** * Protobuf type {@code RequestSpan} @@ -1322,6 +1336,12 @@ public final class TraceProtocol { bussinessKey_ = bs; break; } + case 98: { + com.google.protobuf.ByteString bs = input.readBytes(); + bitField0_ |= 0x00000800; + agentId_ = bs; + break; + } } } } catch (com.google.protobuf.InvalidProtocolBufferException e) { @@ -1743,6 +1763,48 @@ public final class TraceProtocol { } } + public static final int AGENTID_FIELD_NUMBER = 12; + private java.lang.Object agentId_; + /** + * required string agentId = 12; + */ + public boolean hasAgentId() { + return ((bitField0_ & 0x00000800) == 0x00000800); + } + /** + * required string agentId = 12; + */ + public java.lang.String getAgentId() { + java.lang.Object ref = agentId_; + if (ref instanceof java.lang.String) { + return (java.lang.String) ref; + } else { + com.google.protobuf.ByteString bs = + (com.google.protobuf.ByteString) ref; + java.lang.String s = bs.toStringUtf8(); + if (bs.isValidUtf8()) { + agentId_ = s; + } + return s; + } + } + /** + * required string agentId = 12; + */ + public com.google.protobuf.ByteString + getAgentIdBytes() { + java.lang.Object ref = agentId_; + if (ref instanceof java.lang.String) { + com.google.protobuf.ByteString b = + com.google.protobuf.ByteString.copyFromUtf8( + (java.lang.String) ref); + agentId_ = b; + return b; + } else { + return (com.google.protobuf.ByteString) ref; + } + } + private void initFields() { traceId_ = ""; parentLevel_ = ""; @@ -1755,6 +1817,7 @@ public final class TraceProtocol { applicationId_ = ""; userId_ = ""; bussinessKey_ = ""; + agentId_ = ""; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -1798,6 +1861,10 @@ public final class TraceProtocol { memoizedIsInitialized = 0; return false; } + if (!hasAgentId()) { + memoizedIsInitialized = 0; + return false; + } memoizedIsInitialized = 1; return true; } @@ -1838,6 +1905,9 @@ public final class TraceProtocol { if (((bitField0_ & 0x00000400) == 0x00000400)) { output.writeBytes(11, getBussinessKeyBytes()); } + if (((bitField0_ & 0x00000800) == 0x00000800)) { + output.writeBytes(12, getAgentIdBytes()); + } getUnknownFields().writeTo(output); } @@ -1891,6 +1961,10 @@ public final class TraceProtocol { size += com.google.protobuf.CodedOutputStream .computeBytesSize(11, getBussinessKeyBytes()); } + if (((bitField0_ & 0x00000800) == 0x00000800)) { + size += com.google.protobuf.CodedOutputStream + .computeBytesSize(12, getAgentIdBytes()); + } size += getUnknownFields().getSerializedSize(); memoizedSerializedSize = size; return size; @@ -2030,6 +2104,8 @@ public final class TraceProtocol { bitField0_ = (bitField0_ & ~0x00000200); bussinessKey_ = ""; bitField0_ = (bitField0_ & ~0x00000400); + agentId_ = ""; + bitField0_ = (bitField0_ & ~0x00000800); return this; } @@ -2102,6 +2178,10 @@ public final class TraceProtocol { to_bitField0_ |= 0x00000400; } result.bussinessKey_ = bussinessKey_; + if (((from_bitField0_ & 0x00000800) == 0x00000800)) { + to_bitField0_ |= 0x00000800; + } + result.agentId_ = agentId_; result.bitField0_ = to_bitField0_; onBuilt(); return result; @@ -2167,6 +2247,11 @@ public final class TraceProtocol { bussinessKey_ = other.bussinessKey_; onChanged(); } + if (other.hasAgentId()) { + bitField0_ |= 0x00000800; + agentId_ = other.agentId_; + onChanged(); + } this.mergeUnknownFields(other.getUnknownFields()); return this; } @@ -2208,6 +2293,10 @@ public final class TraceProtocol { return false; } + if (!hasAgentId()) { + + return false; + } return true; } @@ -2934,6 +3023,82 @@ public final class TraceProtocol { return this; } + private java.lang.Object agentId_ = ""; + /** + * required string agentId = 12; + */ + public boolean hasAgentId() { + return ((bitField0_ & 0x00000800) == 0x00000800); + } + /** + * required string agentId = 12; + */ + public java.lang.String getAgentId() { + java.lang.Object ref = agentId_; + if (!(ref instanceof java.lang.String)) { + com.google.protobuf.ByteString bs = + (com.google.protobuf.ByteString) ref; + java.lang.String s = bs.toStringUtf8(); + if (bs.isValidUtf8()) { + agentId_ = s; + } + return s; + } else { + return (java.lang.String) ref; + } + } + /** + * required string agentId = 12; + */ + public com.google.protobuf.ByteString + getAgentIdBytes() { + java.lang.Object ref = agentId_; + if (ref instanceof String) { + com.google.protobuf.ByteString b = + com.google.protobuf.ByteString.copyFromUtf8( + (java.lang.String) ref); + agentId_ = b; + return b; + } else { + return (com.google.protobuf.ByteString) ref; + } + } + /** + * required string agentId = 12; + */ + public Builder setAgentId( + java.lang.String value) { + if (value == null) { + throw new NullPointerException(); + } + bitField0_ |= 0x00000800; + agentId_ = value; + onChanged(); + return this; + } + /** + * required string agentId = 12; + */ + public Builder clearAgentId() { + bitField0_ = (bitField0_ & ~0x00000800); + agentId_ = getDefaultInstance().getAgentId(); + onChanged(); + return this; + } + /** + * required string agentId = 12; + */ + public Builder setAgentIdBytes( + com.google.protobuf.ByteString value) { + if (value == null) { + throw new NullPointerException(); + } + bitField0_ |= 0x00000800; + agentId_ = value; + onChanged(); + return this; + } + // @@protoc_insertion_point(builder_scope:RequestSpan) } @@ -2967,14 +3132,14 @@ public final class TraceProtocol { "\n\023TraceProtocol.proto\"z\n\007AckSpan\022\017\n\007trac" + "eId\030\001 \002(\t\022\023\n\013parentLevel\030\002 \001(\t\022\017\n\007levelI" + "d\030\003 \002(\005\022\014\n\004cost\030\004 \002(\003\022\022\n\nstatusCode\030\005 \002(" + - "\005\022\026\n\016exceptionStack\030\006 \001(\t\"\343\001\n\013RequestSpa" + + "\005\022\026\n\016exceptionStack\030\006 \001(\t\"\364\001\n\013RequestSpa" + "n\022\017\n\007traceId\030\001 \002(\t\022\023\n\013parentLevel\030\002 \001(\t\022" + "\017\n\007levelId\030\003 \002(\005\022\023\n\013viewPointId\030\004 \002(\t\022\021\n" + "\tstartDate\030\005 \002(\003\022\024\n\014spanTypeDesc\030\006 \002(\t\022\020" + "\n\010callType\030\007 \002(\t\022\020\n\010spanType\030\010 \002(\r\022\025\n\rap" + "plicationId\030\t \002(\t\022\016\n\006userId\030\n \002(\t\022\024\n\014bus" + - "sinessKey\030\013 \001(\tB(\n&com.ai.cloud.skywalki", - "ng.protocol.proto" + "sinessKey\030\013 \001(\t\022\017\n\007agentId\030\014 \002(\tB(\n&com.", + "ai.cloud.skywalking.protocol.proto" }; com.google.protobuf.Descriptors.FileDescriptor.InternalDescriptorAssigner assigner = new com.google.protobuf.Descriptors.FileDescriptor. InternalDescriptorAssigner() { @@ -2999,7 +3164,7 @@ public final class TraceProtocol { internal_static_RequestSpan_fieldAccessorTable = new com.google.protobuf.GeneratedMessage.FieldAccessorTable( internal_static_RequestSpan_descriptor, - new java.lang.String[] { "TraceId", "ParentLevel", "LevelId", "ViewPointId", "StartDate", "SpanTypeDesc", "CallType", "SpanType", "ApplicationId", "UserId", "BussinessKey", }); + new java.lang.String[] { "TraceId", "ParentLevel", "LevelId", "ViewPointId", "StartDate", "SpanTypeDesc", "CallType", "SpanType", "ApplicationId", "UserId", "BussinessKey", "AgentId", }); } // @@protoc_insertion_point(outer_class_scope) diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java index b28c64abf..ebf55c33e 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java @@ -1,6 +1,7 @@ package com.ai.cloud.skywalking.protocol; import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; import com.ai.cloud.skywalking.protocol.proto.TraceProtocol; import com.google.protobuf.InvalidProtocolBufferException; @@ -24,18 +25,18 @@ public class AckSpan extends AbstractDataSerializable { * 当前调用链的本机描述
* 如当前序号为:0.1.0时,levelId=0 */ - private int levelId = 0; + private int levelId = 0; /** * 节点调用花费时间 */ - private long cost = 0L; + private long cost = 0L; /** * 节点调用的状态
* 0:成功
* 1:异常
* 异常判断原则:代码产生exception,并且此exception不在忽略列表中 */ - private byte statusCode = 0; + private byte statusCode = 0; /** * 节点调用的错误堆栈
* 堆栈以JAVA的exception为主要判断依据 @@ -124,11 +125,12 @@ public class AckSpan extends AbstractDataSerializable { @Override public byte[] getData() { return TraceProtocol.AckSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel). - setLevelId(levelId).setCost(cost).setStatusCode(statusCode).setExceptionStack(exceptionStack).build().toByteArray(); + setLevelId(levelId).setCost(cost).setStatusCode(statusCode).setExceptionStack(exceptionStack).build() + .toByteArray(); } @Override - public AbstractDataSerializable convertData(byte[] data) { + public AbstractDataSerializable convertData(byte[] data) throws ConvertFailedException { AckSpan ackSpan = new AckSpan(); try { TraceProtocol.AckSpan ackSpanProtocol = TraceProtocol.AckSpan.parseFrom(data); @@ -139,7 +141,7 @@ public class AckSpan extends AbstractDataSerializable { ackSpan.setExceptionStack(ackSpanProtocol.getExceptionStack()); ackSpan.setStatusCode((byte) ackSpanProtocol.getStatusCode()); } catch (InvalidProtocolBufferException e) { - return null; + throw new ConvertFailedException(); } return ackSpan; diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java index 38641eae7..5f4b903a9 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java @@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.protocol; import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; import com.ai.cloud.skywalking.protocol.common.CallType; import com.ai.cloud.skywalking.protocol.common.SpanType; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; import com.ai.cloud.skywalking.protocol.proto.TraceProtocol; import com.google.protobuf.InvalidProtocolBufferException; @@ -79,6 +80,11 @@ public class RequestSpan extends AbstractDataSerializable { */ private String businessKey = ""; + /** + * 实例ID + */ + private String agentId = ""; + public RequestSpan(Span spanData) { this.traceId = spanData.getTraceId(); this.parentLevel = spanData.getParentLevel(); @@ -191,13 +197,15 @@ public class RequestSpan extends AbstractDataSerializable { @Override public byte[] getData() { - return TraceProtocol.RequestSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel).setLevelId(levelId).setViewPointId(viewPointId).setStartDate(startDate) - .setSpanType(spanType.getValue()).setSpanTypeDesc(spanTypeDesc).setBussinessKey(businessKey).setCallType(callType).setApplicationId(applicationId).setUserId(userId) - .build().toByteArray(); + return TraceProtocol.RequestSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel) + .setLevelId(levelId).setViewPointId(viewPointId).setStartDate(startDate) + .setSpanType(spanType.getValue()).setSpanTypeDesc(spanTypeDesc).setBussinessKey(businessKey) + .setCallType(callType).setApplicationId(applicationId).setUserId(userId).setBussinessKey(businessKey) + .setAgentId(agentId).build().toByteArray(); } @Override - public AbstractDataSerializable convertData(byte[] data) { + public AbstractDataSerializable convertData(byte[] data) throws ConvertFailedException { RequestSpan requestSpan = new RequestSpan(); try { TraceProtocol.RequestSpan requestSpanByte = TraceProtocol.RequestSpan.parseFrom(data); @@ -211,8 +219,10 @@ public class RequestSpan extends AbstractDataSerializable { requestSpan.setStartDate(requestSpanByte.getStartDate()); requestSpan.setUserId(requestSpanByte.getUserId()); requestSpan.setViewPointId(requestSpanByte.getViewPointId()); + requestSpan.setBusinessKey(requestSpanByte.getBussinessKey()); + requestSpan.setAgentId(requestSpanByte.getAgentId()); } catch (InvalidProtocolBufferException e) { - return null; + throw new ConvertFailedException(); } return requestSpan; @@ -223,6 +233,17 @@ public class RequestSpan extends AbstractDataSerializable { return false; } + public void setBusinessKey(String businessKey) { + this.businessKey = businessKey; + } + + public String getAgentId() { + return agentId; + } + + public void setAgentId(String agentId) { + this.agentId = agentId; + } public static class RequestSpanBuilder { private RequestSpan ackSpan; diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java index 62faaf672..8fd2d7f25 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java @@ -2,6 +2,7 @@ package com.ai.cloud.skywalking.protocol; import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; import com.ai.cloud.skywalking.protocol.common.NullableClass; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; import com.ai.cloud.skywalking.protocol.util.IntegerAssist; import java.util.HashMap; @@ -9,7 +10,8 @@ import java.util.Map; import java.util.ServiceLoader; public class SerializedFactory { - public static Map serializableMap = new HashMap(); + public static Map serializableMap = + new HashMap(); static { @@ -20,13 +22,17 @@ public class SerializedFactory { } } - public static AbstractDataSerializable unSerialize(byte[] bytes) { - AbstractDataSerializable abstractDataSerializable = serializableMap.get(IntegerAssist.bytesToInt(bytes, 0)); - if (abstractDataSerializable != null) { - NullableClass nullableClass = abstractDataSerializable.convert2Object(bytes); - if (!nullableClass.isNull()) { - return (AbstractDataSerializable) nullableClass; + public static AbstractDataSerializable unSerialize(byte[] bytes) throws ConvertFailedException { + try { + AbstractDataSerializable abstractDataSerializable = serializableMap.get(IntegerAssist.bytesToInt(bytes, 0)); + if (abstractDataSerializable != null) { + NullableClass nullableClass = abstractDataSerializable.convert2Object(bytes); + if (!nullableClass.isNull()) { + return (AbstractDataSerializable) nullableClass; + } } + } catch (Exception e) { + throw new ConvertFailedException(); } return null; } diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java index 7c29ccbac..f8b71aa23 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java @@ -49,7 +49,7 @@ public class Span { * 节点的状态
* 不参与序列化 */ - protected boolean isInvalidate = false; + protected boolean isValidate = true; /** * 节点调用过程中的业务字段
@@ -111,8 +111,8 @@ public class Span { this.startDate = startDate; } - public boolean isInvalidate() { - return isInvalidate; + public boolean isValidate() { + return isValidate; } public byte getStatusCode() { @@ -135,8 +135,8 @@ public class Span { this.parameters = parameters; } - public void setInvalidate(boolean invalidate) { - isInvalidate = invalidate; + public void setValidate(boolean validate) { + isValidate = validate; } public boolean isRPCClientSpan() { @@ -154,10 +154,6 @@ public class Span { return spanType; } - public void setIsInvalidate(boolean isInvalidate) { - this.isInvalidate = isInvalidate; - } - public void handleException(Throwable e, Set exclusiveExceptionSet, int maxExceptionStackLength) { ByteArrayOutputStream buf = null; StringBuilder expMessage = new StringBuilder(); @@ -217,4 +213,6 @@ public class Span { public void setUserId(String userId) { this.userId = userId; } + + } diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java index 83e721743..c3eac035f 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java @@ -16,44 +16,30 @@ public class TransportPackager { return dataPackage; } - public static List unpack(byte[] dataPackage) { - if (validateCheckSum(dataPackage)) { - return unpackDataBody(unpackCheckSum(dataPackage)); - } else { - return new ArrayList(); - } - } - - private static byte[] unpackCheckSum(byte[] dataPackage) { - return Arrays.copyOfRange(dataPackage, 4, dataPackage.length); - } - - private static List unpackDataBody(byte[] dataPackage) { - List serializeData = new ArrayList(); - int currentLength = 0; - while (true) { - //读取长度 - int dataLength = IntegerAssist.bytesToInt(dataPackage, currentLength); - // 反序列化 - byte[] data = new byte[dataLength]; - System.arraycopy(dataPackage, currentLength + 4, data, 0, dataLength); - // - serializeData.add(data); - currentLength = 4 + dataLength; - if (currentLength >= dataPackage.length) { - break; + public static List unpackDataBody(byte[] dataPackage) { + List serializeData = null; + try { + serializeData = new ArrayList(); + int currentLength = 0; + while (true) { + //读取长度 + int dataLength = IntegerAssist.bytesToInt(dataPackage, currentLength); + // 反序列化 + byte[] data = new byte[dataLength]; + System.arraycopy(dataPackage, currentLength + 4, data, 0, dataLength); + // + serializeData.add(data); + currentLength += 4 + dataLength; + if (currentLength >= dataPackage.length) { + break; + } } + } catch (Exception e) { + e.printStackTrace(); } return serializeData; } - private static boolean validateCheckSum(byte[] dataPackage) { - byte[] checkSum = generateChecksum(dataPackage, 4); - byte[] originCheckSum = new byte[4]; - System.arraycopy(dataPackage, 0, originCheckSum, 0, 4); - return Arrays.equals(checkSum, originCheckSum); - } - /** * 生成校验和参数 * diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java index 155927931..c573fe20d 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java @@ -2,6 +2,7 @@ package com.ai.cloud.skywalking.protocol.common; import com.ai.cloud.skywalking.protocol.NullClass; import com.ai.cloud.skywalking.protocol.SerializedFactory; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; import com.ai.cloud.skywalking.protocol.util.IntegerAssist; import java.util.Arrays; @@ -14,7 +15,7 @@ public abstract class AbstractDataSerializable implements ISerializable, Nullabl public abstract byte[] getData(); - public abstract AbstractDataSerializable convertData(byte[] data); + public abstract AbstractDataSerializable convertData(byte[] data) throws ConvertFailedException; /** * 消息包结构: @@ -41,7 +42,7 @@ public abstract class AbstractDataSerializable implements ISerializable, Nullabl } @Override - public NullableClass convert2Object(byte[] data) { + public NullableClass convert2Object(byte[] data) throws ConvertFailedException { int dataType = IntegerAssist.bytesToInt(data, 0); if (!SerializedFactory.isCanSerialized(dataType)) { diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/ISerializable.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/ISerializable.java index 81ce0d536..97cfe832e 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/ISerializable.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/ISerializable.java @@ -1,11 +1,13 @@ package com.ai.cloud.skywalking.protocol.common; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; + /** * Created by wusheng on 16/7/4. */ public interface ISerializable { byte[] convert2Bytes(); - NullableClass convert2Object(byte[] data); + NullableClass convert2Object(byte[] data) throws ConvertFailedException; } diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/exception/ConvertFailedException.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/exception/ConvertFailedException.java new file mode 100644 index 000000000..63bd6a44d --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/exception/ConvertFailedException.java @@ -0,0 +1,4 @@ +package com.ai.cloud.skywalking.protocol.exception; + +public class ConvertFailedException extends Exception { +} diff --git a/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto b/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto index 214c3f6f8..77b490abd 100644 --- a/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto +++ b/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto @@ -1,8 +1,8 @@ -syntax="proto2"; +syntax = "proto2"; option java_package = "com.ai.cloud.skywalking.protocol.proto"; -message AckSpan{ +message AckSpan { required string traceId = 1; optional string parentLevel = 2; required int32 levelId = 3; @@ -24,4 +24,5 @@ message RequestSpan { required string applicationId = 9; required string userId = 10; optional string bussinessKey = 11; + required string agentId = 12; } diff --git a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml index c4ff7e408..1e5460750 100644 --- a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml +++ b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml @@ -18,11 +18,6 @@ - - com.ai.cloud - skywalking-api - ${project.version} - com.ai.cloud skywalking-api diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java index ce4187873..cee442c58 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java @@ -4,6 +4,7 @@ import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer; import com.ai.cloud.skywalking.reciever.handler.CollectionServerDataHandler; +import com.ai.cloud.skywalking.reciever.peresistent.PersistenceThreadLauncher; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import io.netty.bootstrap.ServerBootstrap; import io.netty.channel.*; @@ -35,7 +36,8 @@ public class CollectionServer { EventLoopGroup workerGroup = new NioEventLoopGroup(); try { ServerBootstrap b = new ServerBootstrap(); - b.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class).option(ChannelOption.SO_BACKLOG, 100).handler(new LoggingHandler(LogLevel.INFO)) + b.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class).option(ChannelOption.SO_BACKLOG, 100) + .handler(new LoggingHandler(LogLevel.INFO)) .childHandler(new ChannelInitializer() { @Override public void initChannel(io.netty.channel.socket.SocketChannel ch) throws Exception { @@ -61,6 +63,8 @@ public class CollectionServer { initializeParam(); logger.info("To init server health collector..."); ServerHealthCollector.init(); + logger.info("To launch register persistence thread...."); + PersistenceThreadLauncher.doLaunch(); logger.info("To init data buffer thread container..."); DataBufferThreadContainer.init(); logger.info("Starting collection server....."); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java new file mode 100644 index 000000000..c9a6d1020 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java @@ -0,0 +1,44 @@ +package com.ai.cloud.skywalking.reciever.buffer; + +import com.ai.cloud.skywalking.reciever.model.BufferDataPackagerGenerator; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.io.File; +import java.io.FileOutputStream; +import java.io.IOException; +import java.util.concurrent.CountDownLatch; + +class AppendEOFFlagThread extends Thread { + private Logger logger = LogManager.getLogger(AppendEOFFlagThread.class); + private File[] dataBufferFiles; + private CountDownLatch countDownLatch; + + public AppendEOFFlagThread(File[] dataBufferFiles, CountDownLatch countDownLatch) { + super("AppendEOFFlagThread"); + this.dataBufferFiles = dataBufferFiles; + this.countDownLatch = countDownLatch; + } + + @Override + public void run() { + FileOutputStream fileOutputStream = null; + for (File file : dataBufferFiles) { + try { + logger.info("Add EOF flags to unprocessed data file[{}]", file.getName()); + fileOutputStream = new FileOutputStream(new File(file.getParent(), file.getName()), true); + fileOutputStream.write(BufferDataPackagerGenerator.generateEOFPackage()); + } catch (IOException e) { + logger.info("Add EOF flags to the unprocessed data file failed.", e); + } finally { + try { + fileOutputStream.flush(); + fileOutputStream.close(); + } catch (IOException e) { + logger.error("Flush data file failed", e); + } + } + countDownLatch.countDown(); + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java index 05bde4ee7..26f41e02b 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java @@ -1,20 +1,18 @@ package com.ai.cloud.skywalking.reciever.buffer; -import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; +import com.ai.cloud.skywalking.protocol.util.AtomicRangeInteger; import com.ai.cloud.skywalking.reciever.conf.Config; -import com.ai.cloud.skywalking.reciever.processor.IProcessor; -import com.ai.cloud.skywalking.reciever.processor.ProcessorFactory; +import com.ai.cloud.skywalking.reciever.model.BufferDataPackagerGenerator; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; -import com.ai.cloud.skywalking.protocol.SerializedFactory; -import com.ai.cloud.skywalking.protocol.util.AtomicRangeInteger; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.io.File; +import java.io.FileNotFoundException; +import java.io.FileOutputStream; +import java.io.IOException; +import java.util.UUID; import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.DATA_CONFLICT_WAIT_TIME; import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.PER_THREAD_MAX_BUFFER_NUMBER; @@ -31,28 +29,30 @@ public class DataBufferThread extends Thread { @Override public void run() { - Map> serializeObjects; + FileOutputStream fileOutputStream = null; + int length = 0; while (true) { - serializeObjects = new HashMap>(); for (int i = 0; i < data.length; i++) { if (data[i] == null) { continue; } - AbstractDataSerializable serializeData = SerializedFactory.unSerialize(data[i]); - List hasBeenSerializedObjects = serializeObjects.get(serializeData.getDataType()); - if (hasBeenSerializedObjects == null) { - serializeObjects.put(serializeData.getDataType(), new ArrayList()); + if (fileOutputStream == null) { + fileOutputStream = acquiredNewBufferFileStream(); } - serializeObjects.get(serializeData.getDataType()).add(serializeData); - data[i] = null; - } + try { + fileOutputStream.write(BufferDataPackagerGenerator.pack(data[i])); + length += data[i].length; + data[i] = null; + } catch (IOException e) { + logger.error("Failed to write msg.", e); + } - for (Map.Entry> entry : serializeObjects.entrySet()) { - IProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey()); - if (processor != null) { - processor.process(entry.getValue()); + + if (length > Config.Buffer.BUFFER_FILE_MAX_LENGTH) { + closeCurrentBufferFile(fileOutputStream); + fileOutputStream = null; } } @@ -65,19 +65,73 @@ public class DataBufferThread extends Thread { } } + private void closeCurrentBufferFile(FileOutputStream fileOutputStream) { + try { + fileOutputStream.flush(); + fileOutputStream.write(BufferDataPackagerGenerator.generateEOFPackage()); + } catch (IOException e) { + logger.error("Failed to write msg.", e); + } finally { + try { + fileOutputStream.flush(); + fileOutputStream.close(); + } catch (IOException e) { + logger.error("Failed to flush and close file.", e); + } + } + } + private FileOutputStream acquiredNewBufferFileStream() { + checkBufferDirIsExists(); + File outputFile = createNewBufferFile(); + return generateFileOutputStream(outputFile); + } + + private FileOutputStream generateFileOutputStream(File outputFile) { + FileOutputStream fileOutputStream = null; + try { + fileOutputStream = new FileOutputStream(outputFile); + } catch (FileNotFoundException e) { + logger.error("Failed to create File:{}", outputFile.getName(), e); + } + return fileOutputStream; + } + + private File createNewBufferFile() { + File outputFile = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR, + System.currentTimeMillis() + "-" + UUID.randomUUID().toString()); + if (!outputFile.exists()) { + try { + outputFile.createNewFile(); + ServerHealthCollector.getCurrentHeathReading(null) + .updateData(ServerHeathReading.INFO, "Create new Buffer file[" + outputFile.getName() + "]"); + } catch (IOException e) { + logger.error("Failed to create File:{}", outputFile.getName(), e); + } + } + return outputFile; + } + + private void checkBufferDirIsExists() { + File outPutDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR); + if (!outPutDir.exists()) { + outPutDir.mkdirs(); + } + } public void saveTemporarily(byte[] s) { int i = index.getAndIncrement(); while (data[i] != null) { try { - ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.WARNING, "DataBuffer index[" + i + "] data collision, service pausing. "); + ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.WARNING, + "DataBuffer index[" + i + "] data collision, service pausing. "); Thread.sleep(DATA_CONFLICT_WAIT_TIME); } catch (InterruptedException e) { logger.error("Failure sleep.", e); } } - ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.INFO, "DataBuffer reveiving data."); + ServerHealthCollector.getCurrentHeathReading(null) + .updateData(ServerHeathReading.INFO, "DataBuffer reveiving data."); data[i] = s; } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java index 98764a344..f40d4dd94 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java @@ -1,13 +1,19 @@ package com.ai.cloud.skywalking.reciever.buffer; import com.ai.cloud.skywalking.reciever.conf.Config; +import org.apache.commons.io.comparator.NameFileComparator; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import java.io.File; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ThreadLocalRandom; +import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.MAX_APPEND_EOF_FLAGS_THREAD_NUMBER; + public class DataBufferThreadContainer { private static Logger logger = LogManager.getLogger(DataBufferThreadContainer.class); @@ -23,7 +29,38 @@ public class DataBufferThreadContainer { return buffers.get(ThreadLocalRandom.current().nextInt(buffers.size())); } - public static void init() { + public static void init() throws InterruptedException { + logger.info("Add EOF flags to the unprocessed data file last time."); + File parentDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR); + // 判断数据缓存文件是否存在,如果不存在,则创建 + if (!parentDir.exists()) { + parentDir.mkdirs(); + } + NameFileComparator sizeComparator = new NameFileComparator(); + File[] dataFileList = sizeComparator.sort(parentDir.listFiles()); + logger.info("Pending file number :" + dataFileList.length); + if (dataFileList.length > 0) { + int step = (int) Math.ceil(dataFileList.length * 1.0 / MAX_APPEND_EOF_FLAGS_THREAD_NUMBER); + + int start = 0, end = 0; + CountDownLatch countDownLatch = new CountDownLatch(MAX_APPEND_EOF_FLAGS_THREAD_NUMBER); + while (true) { + if (end + step >= dataFileList.length) { + new AppendEOFFlagThread(Arrays.copyOfRange(dataFileList, start, dataFileList.length), + countDownLatch).start(); + break; + } + end += step; + new AppendEOFFlagThread(Arrays.copyOfRange(dataFileList, start, end), countDownLatch).start(); + start = end; + logger.debug("start:" + start + "\tend:" + end); + } + + countDownLatch.await(); + } + logger.info("Data buffer thread size {} begin to init ", Config.Server. + MAX_DEAL_DATA_THREAD_NUMBER); + for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { DataBufferThread dataBufferThread = new DataBufferThread(i); dataBufferThread.start(); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java index 234310c76..f2cb5d942 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java @@ -17,6 +17,8 @@ public class Config { // 数据缓存配置类 public static class Buffer { + // 数据缓存文件目录 + public static String DATA_BUFFER_FILE_PARENT_DIR = "/tmp/skywalking/data/buffer"; //每个线程最大缓存数量 public static int PER_THREAD_MAX_BUFFER_NUMBER = 1024; @@ -27,6 +29,33 @@ public class Config { // 数据冲突时等待时间(单位:毫秒) public static long DATA_CONFLICT_WAIT_TIME = 10L; + public static long BUFFER_FILE_MAX_LENGTH = 30 * 1024 * 1024; + + } + + + public static class RegisterPersistence { + // 偏移量注册文件的目录 + public static String REGISTER_FILE_PARENT_DIRECTORY = "../data/offset"; + + // 偏移量注册文件名 + public static String REGISTER_FILE_NAME = "offset.A"; + + // 偏移量注册备份文件名 + public static String REGISTER_BAK_FILE_NAME = "offset.B"; + + // 偏移量写入文件等待周期 + public static long OFFSET_WRITTEN_FILE_WAIT_CYCLE = 5000L; + } + + + public static class Persistence { + // 切换文件,等待时间 + public static long SWITCH_FILE_WAIT_TIME = 5000L; + + // 追加EOF标志位的线程数量 + public static int MAX_APPEND_EOF_FLAGS_THREAD_NUMBER = 1; + } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java index 9bd89925f..958b5b095 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java @@ -1,15 +1,15 @@ package com.ai.cloud.skywalking.reciever.handler; +import com.ai.cloud.skywalking.protocol.util.IntegerAssist; import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.util.RedisConnector; -import com.ai.cloud.skywalking.protocol.TransportPackager; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import redis.clients.jedis.Jedis; import java.net.InetSocketAddress; -import java.util.List; +import java.util.Arrays; public class CollectionServerDataHandler extends SimpleChannelInboundHandler { @@ -18,11 +18,8 @@ public class CollectionServerDataHandler extends SimpleChannelInboundHandler= 0) { - - List byteSerializeObjects = TransportPackager.unpack(msg); - - if (byteSerializeObjects.size() > 0) { - cacheSerializeObjects(byteSerializeObjects); + if (validateCheckSum(msg)) { + DataBufferThreadContainer.getDataBufferThread().saveTemporarily(unpackCheckSum(msg)); } else { // 处理错误包 dealFailedPackage(ctx); @@ -30,12 +27,27 @@ public class CollectionServerDataHandler extends SimpleChannelInboundHandler byteSerializeObjects) { - for (byte[] byteSerializeObject : byteSerializeObjects) { - DataBufferThreadContainer.getDataBufferThread().saveTemporarily(byteSerializeObject); - } + private byte[] unpackCheckSum(byte[] msg) { + return Arrays.copyOfRange(msg, 4, msg.length); } + private boolean validateCheckSum(byte[] dataPackage) { + byte[] checkSum = generateChecksum(dataPackage, 4); + byte[] originCheckSum = new byte[4]; + System.arraycopy(dataPackage, 0, originCheckSum, 0, 4); + return Arrays.equals(checkSum, originCheckSum); + } + + private static byte[] generateChecksum(byte[] data, int offset) { + int result = data[offset]; + for (int i = offset + 1; i < data.length; i++) { + result ^= data[i]; + } + + return IntegerAssist.intToBytes(result); + } + + private void dealFailedPackage(ChannelHandlerContext ctx) { InetSocketAddress socketAddress = (InetSocketAddress) ctx.channel().localAddress(); String key = ctx.name() + "-" + socketAddress.getHostName() + ":" + socketAddress.getPort(); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java new file mode 100644 index 000000000..a2ff30dd0 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java @@ -0,0 +1,39 @@ +package com.ai.cloud.skywalking.reciever.model; + +import com.ai.cloud.skywalking.protocol.util.IntegerAssist; + +public class BufferDataPackagerGenerator { + + private static byte[] SPILT = new byte[] {127, 127, 127, 127}; + private static byte[] EOF = null; + + private BufferDataPackagerGenerator() { + //DO Nothing + } + + public static byte[] generateEOFPackage() { + if (EOF != null) { + return EOF; + } + + EOF = generatePackage("EOF".getBytes()); + return EOF; + } + + public static byte[] pack(byte[] msg) { + return generatePackage(msg); + } + + + private static byte[] generatePackage(byte[] msg) { + byte[] dataPackage = new byte[msg.length + 8]; + // 前四位长度 + System.arraycopy(IntegerAssist.intToBytes(msg.length), 0, dataPackage, 0, 4); + // 中间正文 + System.arraycopy(msg, 0, dataPackage, 4, msg.length); + // 后四位特殊字符 + System.arraycopy(SPILT, 0, dataPackage, msg.length + 4, 4); + + return dataPackage; + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java new file mode 100644 index 000000000..e5a378da6 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java @@ -0,0 +1,227 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.protocol.util.IntegerAssist; +import org.apache.commons.io.FileUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.io.File; +import java.io.FileInputStream; +import java.io.IOException; +import java.util.Arrays; + +public class BufferBalePlucker { + private File bufferFile; + private FileInputStream bufferInputStream; + private int currentOffset; + private boolean hasNextBufferBale = false; + private static final byte[] SPILT_BALE_ARRAY = new byte[] {127, 127, 127, 127}; + private static final byte[] EOF_BALE_ARRAY = "EOF".getBytes(); + private int remainderLength = 0; + private byte[] remainderByte = null; + private Logger logger = LogManager.getLogger(BufferBalePlucker.class); + private PluckerStatus status = PluckerStatus.INITIAL; + + public BufferBalePlucker(File bufferFile, int currentOffset) { + this.bufferFile = bufferFile; + this.currentOffset = currentOffset; + + if (bufferFile.length() >= currentOffset) { + hasNextBufferBale = true; + } + + try { + this.bufferInputStream = new FileInputStream(bufferFile); + bufferInputStream.skip(currentOffset); + } catch (IOException e) { + hasNextBufferBale = false; + } + + } + + public boolean hasNextBufferBale() { + + if (status == PluckerStatus.SUSPEND) { + try { + Thread.sleep(500); + } catch (InterruptedException e) { + logger.error("Failed to sleep ", e); + } + } + + if (status == PluckerStatus.INITIAL) { + status = PluckerStatus.RUNNING; + return hasNextBufferBale; + } + + if (status == PluckerStatus.RUNNING) { + byte[] spiltArray = new byte[4]; + try { + bufferInputStream.read(spiltArray); + if (!Arrays.equals(spiltArray, SPILT_BALE_ARRAY)) { + skipToNextBufferBale(); + } + } catch (IOException e) { + return false; + } + currentOffset += 4; + } + MemoryRegister.instance().updateOffSet(bufferFile.getName(), currentOffset); + return hasNextBufferBale; + } + + public byte[] pluck() throws IOException { + int packageLength = unpackBaleLength(); + byte[] dataPackage = unpackDataContext(packageLength); + + if (dataPackage == null || dataPackage.length == 0) { + // 文件没有完结,但是已经被处理完成了 + return null; + } + + if (checkDataPackageIsEOF(dataPackage)) { + hasNextBufferBale = false; + return null; + } + + if (logger.isDebugEnabled()) { + logger.debug("Pluck bale size : " + dataPackage.length); + } + + return dataPackage; + } + + private boolean checkDataPackageIsEOF(byte[] dataPackage) { + if (dataPackage.length == EOF_BALE_ARRAY.length) { + return Arrays.equals(dataPackage, EOF_BALE_ARRAY); + } + return false; + } + + private byte[] unpackDataContext(int length) throws IOException { + if (currentOffset >= bufferFile.length()) { + status = PluckerStatus.SUSPEND; + return null; + } + byte[] dataContext = new byte[length]; + bufferInputStream.read(dataContext); + currentOffset += length; + return dataContext; + } + + private int unpackBaleLength() { + int length; + + while (true) { + try { + length = calculateCurrentPackageLength(); + if (length > 0 && length < 90000) { + break; + } + skipToNextBufferBale(); + } catch (IOException e) { + skipToNextBufferBale(); + } + } + + return length; + } + + private int calculateCurrentPackageLength() throws IOException { + byte[] lengthByte = new byte[4 - remainderLength]; + bufferInputStream.read(lengthByte); + currentOffset += 4 - remainderLength; + lengthByte = spliceRemainderByteOfPreviousSkipIfNecessary(lengthByte); + return IntegerAssist.bytesToInt(lengthByte, 0); + } + + private byte[] spliceRemainderByteOfPreviousSkipIfNecessary(byte[] lengthByte) { + if (remainderLength != 0) { + byte[] length = Arrays.copyOf(remainderByte, 4); + System.arraycopy(lengthByte, 0, length, remainderLength, lengthByte.length); + remainderLength = 0; + return length; + } + return lengthByte; + } + + public void skipToNextBufferBale() { + byte[] previousDataByte = new byte[4]; + byte[] currentDataByte = new byte[4]; + byte[] compactDataByte = new byte[8]; + while (true) { + try { + currentOffset += bufferInputStream.read(currentDataByte); + } catch (IOException e) { + hasNextBufferBale = false; + } + + if (Arrays.equals(currentDataByte, SPILT_BALE_ARRAY)) { + remainderLength = 0; + break; + } + + // + if (currentOffset + 8000 >= bufferFile.length()) { + status = PluckerStatus.SUSPEND; + break; + } + + System.arraycopy(previousDataByte, 0, compactDataByte, 0, 4); + System.arraycopy(currentDataByte, 0, compactDataByte, 4, 4); + + int index = bytesIndexOf(compactDataByte, SPILT_BALE_ARRAY, 0, 8); + if (index != -1) { + recodeRemainderByteAndLength(compactDataByte, index); + break; + } + + previousDataByte = Arrays.copyOf(currentDataByte, 4); + } + } + + private void recodeRemainderByteAndLength(byte[] compactDataByte, int index) { + remainderLength = 8 - 4 - index; + remainderByte = Arrays.copyOfRange(compactDataByte, index + 4, index + 4 + remainderLength); + } + + private int bytesIndexOf(byte[] Source, byte[] Search, int fromIndex, int endIndex) { + boolean find = false; + int i; + for (i = fromIndex; i < endIndex - Search.length; i++) { + if (Source[i] == Search[0]) { + find = true; + for (int j = 0; j < Search.length; j++) { + if (Source[i + j] != Search[j]) { + find = false; + } + } + } + if (find) { + break; + } + } + return !find ? -1 : i; + } + + public void close() throws IOException { + if (bufferInputStream != null) { + bufferInputStream.close(); + } + logger.info("Data in file[{}] has been successfully processed", bufferFile.getName()); + boolean deleteSuccess = false; + while (!deleteSuccess) { + deleteSuccess = FileUtils.deleteQuietly(new File(bufferFile.getParent(), bufferFile.getName())); + } + logger.info("Delete file[{}] {}", bufferFile.getName(), (deleteSuccess ? "success" : "failed")); + MemoryRegister.instance().removeEntry(bufferFile.getName()); + bufferFile = null; + } + + enum PluckerStatus { + INITIAL, + RUNNING, + SUSPEND + } + +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/FileRegisterEntry.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/FileRegisterEntry.java new file mode 100644 index 000000000..a284ddc2c --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/FileRegisterEntry.java @@ -0,0 +1,78 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +public class FileRegisterEntry { + private String fileName; + private int offset; + private FileRegisterEntryStatus status; + + public FileRegisterEntry() { + } + + public FileRegisterEntry(String fileName) { + this.fileName = fileName; + } + + public FileRegisterEntry(String fileName, int offset) { + this.fileName = fileName; + this.offset = offset; + } + + public FileRegisterEntry(String fileName, int offset, FileRegisterEntryStatus status) { + this.fileName = fileName; + this.offset = offset; + this.status = status; + } + + public String getFileName() { + return fileName; + } + + public void setFileName(String fileName) { + this.fileName = fileName; + } + + public int getOffset() { + return offset; + } + + public void setOffset(int offset) { + this.offset = offset; + } + + public FileRegisterEntryStatus getStatus() { + return status; + } + + public void setStatus(FileRegisterEntryStatus status) { + this.status = status; + } + + @Override + public boolean equals(Object o) { + if (this == o) + return true; + if (!(o instanceof FileRegisterEntry)) + return false; + + FileRegisterEntry that = (FileRegisterEntry) o; + + return !(getFileName() != null ? !getFileName().equals(that.getFileName()) : that.getFileName() != null); + + } + + @Override + public int hashCode() { + return getFileName() != null ? getFileName().hashCode() : 0; + } + + + @Override + public String toString() { + return fileName + '\t' + offset + "\t" + status; + } + + public enum FileRegisterEntryStatus { + REGISTER, + UNREGISTER; + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/MemoryRegister.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/MemoryRegister.java new file mode 100644 index 000000000..a5ee75cb5 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/MemoryRegister.java @@ -0,0 +1,127 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.reciever.conf.Config; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.io.File; +import java.io.FileNotFoundException; +import java.io.IOException; +import java.util.Arrays; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.*; + +public class MemoryRegister { + private Logger logger = LogManager.getLogger(MemoryRegister.class); + private Map entries = new ConcurrentHashMap(); + private File file; + private static MemoryRegister memoryRegister = new MemoryRegister(); + + public static MemoryRegister instance() { + return memoryRegister; + } + + public void updateOffSet(String fileName, int offset) { + if (logger.isDebugEnabled()) { + logger.debug("Register entry[{}] offset[{}] into the memory register", fileName, offset); + } + if (entries.containsKey(fileName)) { + entries.get(fileName).setOffset(offset); + } + } + + public void unRegister(String fileName) { + if (logger.isDebugEnabled()) { + logger.debug("Unregister[{}] from the memory register", fileName); + } + if (entries.containsKey(fileName)) { + entries.get(fileName).setStatus(FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER); + } + } + + + public void removeEntry(String fileName) { + entries.remove(fileName); + } + + public synchronized FileRegisterEntry doRegister(String fileName) { + logger.debug("Begin to register File[{}]", fileName); + FileRegisterEntry entry = null; + // 已经存在entries. + if (entries.containsKey(fileName)) { + logger.debug("FileRegisterEntry[{}] Status:[{}]", entries.get(fileName).getStatus()); + // 已经被别的线程处理中 + if (entries.get(fileName).getStatus() == FileRegisterEntry.FileRegisterEntryStatus.REGISTER) { + if (logger.isDebugEnabled()) { + logger.debug("Entry[{}] has been register", fileName); + } + + } else { + // 没有被别的线程处理 + entry = entries.get(fileName); + entry.setStatus(FileRegisterEntry.FileRegisterEntryStatus.REGISTER); + } + } else { + // 以前没有被注册过的 + entry = new FileRegisterEntry(fileName, 0, FileRegisterEntry.FileRegisterEntryStatus.REGISTER); + entries.put(fileName, entry); + } + + return entry; + } + + + public Collection getEntries() { + return entries.values(); + } + + public int getOffSet(String fileName) { + if (entries.containsKey(fileName)) { + return entries.get(fileName).getOffset(); + } + return -1; + } + + private void checkOffSetExists() { + file = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_FILE_NAME); + + if (!file.exists()) { + try { + file.createNewFile(); + } catch (IOException e) { + logger.error("Create offset filed failed", e); + } + } + } + + private MemoryRegister() { + // 在处理数据之前需要初始化处理文件的处理状态 + //去掉entries中无法与缓存数据文件匹配的文件 + File parentDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR); + //上次未处理的缓存数据文件,entries内的数据主要以缓存 + List bufferFileNameList = Arrays.asList(parentDir.list()); + + try { + // 读取offset文件 + OffsetFile mainOffsetFile = new OffsetFile(REGISTER_FILE_PARENT_DIRECTORY + File.separator + + REGISTER_FILE_NAME, bufferFileNameList); + OffsetFile backUpOffsetFile = new OffsetFile(REGISTER_FILE_PARENT_DIRECTORY + File.separator + + REGISTER_BAK_FILE_NAME, bufferFileNameList); + + if (mainOffsetFile.compare(backUpOffsetFile)) { + entries.putAll(mainOffsetFile.getRegisterEntries()); + } else { + entries.putAll(backUpOffsetFile.getRegisterEntries()); + } + } catch (FileNotFoundException e) { + logger.error("The offset file does not exist.", e); + checkOffSetExists(); + } catch (IOException e) { + logger.error("Read data from offset file failed.", e); + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/OffsetFile.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/OffsetFile.java new file mode 100644 index 000000000..546aa0f8f --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/OffsetFile.java @@ -0,0 +1,69 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import java.io.BufferedReader; +import java.io.FileReader; +import java.io.IOException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +public class OffsetFile { + private boolean isComplete = true; + private Map registerEntries = new HashMap(); + private long lastModifyTime; + + public OffsetFile(String offsetFileName, List bufferFileNameList) throws IOException { + BufferedReader reader = new BufferedReader(new FileReader(offsetFileName)); + String offsetData; + String lastModifyTimeStr = reader.readLine(); + if (lastModifyTimeStr == null || lastModifyTimeStr.length() == 0) { + isComplete = false; + return; + } + lastModifyTime = Long.parseLong(lastModifyTimeStr); + while ((offsetData = reader.readLine()) != null && !"EOF".equals(offsetData)) { + String[] ss = offsetData.split("\t"); + if (bufferFileNameList.contains(ss[0])) { + registerEntries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), + FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + } + } + + if (!"EOF".equals(offsetData)) { + // 文件不完整 + isComplete = false; + } + } + + public boolean compare(OffsetFile offsetFile) { + if (isComplete && !offsetFile.isComplete) { + if (lastModifyTime > offsetFile.lastModifyTime) { + // 优先选择完整,并且时间是最新的 + return true; + } else { + return false; + } + } else if (isComplete && offsetFile.isComplete) { + // 都是完整,则采用时间最新的 + if (lastModifyTime > offsetFile.lastModifyTime) { + return true; + } else { + return false; + } + } else if (!isComplete && offsetFile.isComplete) { + //本身没有完成,但是别人完成了,采用别人的 + return false; + } else { + // 自己没有完成,别人也没有完成,则比较时间 + if (lastModifyTime > offsetFile.lastModifyTime) { + return true; + } else { + return false; + } + } + } + + public Map getRegisterEntries() { + return registerEntries; + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java new file mode 100644 index 000000000..7efbe2e60 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java @@ -0,0 +1,142 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.protocol.SerializedFactory; +import com.ai.cloud.skywalking.protocol.TransportPackager; +import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; +import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.processor.IProcessor; +import com.ai.cloud.skywalking.reciever.processor.ProcessorFactory; +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; +import org.apache.commons.io.comparator.NameFileComparator; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.io.File; +import java.io.IOException; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +public class PersistenceThread extends Thread { + + private Logger logger = LogManager.getLogger(PersistenceThread.class); + private File bufferFile = null; + + public PersistenceThread(int trdIndex) { + super("PersistentThread" + trdIndex); + } + + @Override + public void run() { + while (true) { + bufferFile = chooseDealBufferFile(); + + if (bufferFile == null) { + try { + Thread.sleep(Config.Persistence.SWITCH_FILE_WAIT_TIME); + } catch (InterruptedException e) { + logger.error("Failure sleep.", e); + } + continue; + } + + int offset = acquireOffset(); + + BufferBalePlucker bufferBalePlucker = new BufferBalePlucker(bufferFile, offset); + while (bufferBalePlucker.hasNextBufferBale()) { + try { + Map> spans = unSerializeSpans(bufferBalePlucker.pluck()); + System.out.println(spans.size()); + //handleSpans(spans); + } catch (ConvertFailedException e) { + bufferBalePlucker.skipToNextBufferBale(); + } catch (IOException e) { + logger.error("The data file I/O exception.", e); + ServerHealthCollector.getCurrentHeathReading(null) + .updateData(ServerHeathReading.ERROR, e.getMessage()); + } + } + + try { + bufferBalePlucker.close(); + } catch (IOException e) { + logger.error("The data file I/O exception.", e); + ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.ERROR, e.getMessage()); + } + + try { + Thread.sleep(Config.Persistence.SWITCH_FILE_WAIT_TIME); + } catch (InterruptedException e) { + logger.error("Failure sleep.", e); + } + } + } + + private int acquireOffset() { + int offset; + offset = MemoryRegister.instance().getOffSet(bufferFile.getName()); + if (offset == -1 || offset == 0) { + offset = 0; + } + return offset; + } + + private void handleSpans(Map> spans) { + for (Map.Entry> entry : spans.entrySet()) { + IProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey()); + if (processor != null) { + processor.process(entry.getValue()); + } + } + } + + private Map> unSerializeSpans(byte[] byteData) + throws ConvertFailedException { + Map> spans = new HashMap>(); + if (byteData == null || byteData.length == 0) { + return spans; + } + List serializeData = TransportPackager.unpackDataBody(byteData); + for (byte[] dataBytes : serializeData) { + AbstractDataSerializable abstractDataSerializable = SerializedFactory.unSerialize(dataBytes); + if (spans.get(abstractDataSerializable.getDataType()) == null) { + spans.put(abstractDataSerializable.getDataType(), new ArrayList()); + } + spans.get(abstractDataSerializable.getDataType()).add(abstractDataSerializable); + } + + return spans; + } + + private File chooseDealBufferFile() { + File file1 = null; + File parentDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR); + NameFileComparator sizeComparator = new NameFileComparator(); + File[] dataFileList = sizeComparator.sort(parentDir.listFiles()); + + if (dataFileList == null) { + return null; + } + + for (File file : dataFileList) { + if (file.getName().startsWith(".")) { + continue; + } + if (MemoryRegister.instance().doRegister(file.getName()) == null) { + if (logger.isDebugEnabled()) + logger.debug("The file [{}] is being used by another thread ", file); + continue; + } + if (logger.isDebugEnabled()) { + logger.debug("Begin to deal data file [{}]", file.getName()); + } + file1 = file; + break; + } + + return file1; + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java new file mode 100644 index 000000000..8f3d2e526 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java @@ -0,0 +1,12 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.reciever.conf.Config; + +public class PersistenceThreadLauncher { + public static void doLaunch() { + new RegisterPersistenceThread().start(); + for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { + new PersistenceThread(i).start(); + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/RegisterPersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/RegisterPersistenceThread.java new file mode 100644 index 000000000..ab2e7fc87 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/RegisterPersistenceThread.java @@ -0,0 +1,98 @@ +package com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.io.BufferedWriter; +import java.io.File; +import java.io.FileWriter; +import java.io.IOException; +import java.util.Collection; + +import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.*; + +public class RegisterPersistenceThread extends Thread { + private Logger logger = LogManager.getLogger(RegisterPersistenceThread.class); + private File offsetFile; + private File bakOffsetFile; + private BufferedWriter offsetBufferWriter; + private boolean isSwitch = true; + + public RegisterPersistenceThread() { + super("RegisterPersistenceThread"); + File offsetParentDir = new File(REGISTER_FILE_PARENT_DIRECTORY); + if (!offsetParentDir.exists()) { + offsetParentDir.mkdirs(); + } + + offsetFile = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_FILE_NAME); + bakOffsetFile = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_BAK_FILE_NAME); + } + + @Override + public void run() { + + while (true) { + try { + Thread.sleep(OFFSET_WRITTEN_FILE_WAIT_CYCLE); + } catch (InterruptedException e) { + logger.error("Sleep failure", e); + } + chooseOffsetFile(); + appendingTimestamp(); + wirteFileRegisterEntries(); + closeAndReleaseResource(); + ServerHealthCollector.getCurrentHeathReading(null) + .updateData(ServerHeathReading.INFO, "flush memory register to file."); + } + } + + private void chooseOffsetFile() { + try { + if (isSwitch) { + offsetBufferWriter = new BufferedWriter(new FileWriter(offsetFile)); + } else { + offsetBufferWriter = new BufferedWriter(new FileWriter(bakOffsetFile)); + } + isSwitch = !isSwitch; + } catch (IOException e) { + logger.error("Write The offset file anomalies."); + } + } + + private void appendingTimestamp() { + try { + offsetBufferWriter.write(String.valueOf(System.currentTimeMillis()) + "\n"); + } catch (IOException e) { + logger.error("Write The offset file anomalies."); + } + } + + private void wirteFileRegisterEntries() { + Collection fileRegisterEntries = MemoryRegister.instance().getEntries(); + for (FileRegisterEntry fileRegisterEntry : fileRegisterEntries) { + try { + offsetBufferWriter.write(fileRegisterEntry.toString() + "\n"); + } catch (IOException e) { + logger.error("Write file register entry to offset file failure", e); + } + } + } + + private void closeAndReleaseResource() { + try { + offsetBufferWriter.write("EOF\n"); + offsetBufferWriter.flush(); + } catch (IOException e) { + logger.error("Flush offset file failure", e); + } finally { + try { + offsetBufferWriter.close(); + } catch (IOException e) { + logger.error("close offset file failure", e); + } + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java index 14c78c715..68aa284db 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java @@ -14,6 +14,7 @@ import java.util.List; import static com.ai.cloud.skywalking.reciever.util.SpanUtil.getTSBySpanTraceId; +@DefaultProcessor public class AckSpanProcessor extends AbstractSpanProcessor { @Override diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/DefaultProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/DefaultProcessor.java new file mode 100644 index 000000000..e68e5d933 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/DefaultProcessor.java @@ -0,0 +1,12 @@ +package com.ai.cloud.skywalking.reciever.processor; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +@Target(ElementType.TYPE) +@Retention(RetentionPolicy.RUNTIME) +public @interface DefaultProcessor { + boolean defaultProcessor() default true; +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java index e1694a12c..d48eff01b 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java @@ -8,15 +8,19 @@ import java.util.Map; import java.util.ServiceLoader; public class ProcessorFactory { - private static Logger logger = LogManager.getLogger(ProcessorFactory.class); + private static Logger logger = LogManager.getLogger(ProcessorFactory.class); private static Map type_processor_mapping = new HashMap(); static { ServiceLoader processors = ServiceLoader.load(IProcessor.class); - for (IProcessor processor : processors) { - logger.info("Init protocol type and processor mapping : {} --> {}.", processor.getProtocolType(), processor.getClass().getName()); - type_processor_mapping.put(processor.getProtocolType(), processor); + DefaultProcessor defaultProcessor = processor.getClass().getAnnotation(DefaultProcessor.class); + IProcessor processor1 = type_processor_mapping.get(processor.getProtocolType()); + if (processor1 == null || (defaultProcessor != null && !defaultProcessor.defaultProcessor())) { + logger.info("Init protocol type and processor mapping : {} --> {}.", processor.getProtocolType(), + processor.getClass().getName()); + type_processor_mapping.put(processor.getProtocolType(), processor); + } } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java index 715186a37..e46aa9509 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java @@ -14,6 +14,7 @@ import java.util.List; import static com.ai.cloud.skywalking.reciever.util.SpanUtil.getTSBySpanTraceId; +@DefaultProcessor public class RequestSpanProcessor extends AbstractSpanProcessor { @@ -35,7 +36,8 @@ public class RequestSpanProcessor extends AbstractSpanProcessor { } else { columnName = requestSpan.getParentLevel() + "." + requestSpan.getLevelId(); } - put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), requestSpan.getData()); + put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), + requestSpan.getData()); puts.add(put); } diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index a525da474..76995a6e2 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -11,16 +11,26 @@ buffer.max_wait_time=5000 #数据冲突时等待时间(单位:毫秒) buffer.data_conflict_wait_time=10 #数据缓存文件目录 -buffer.data_buffer_file_parent_directory=D:/test-data/data/buffer +buffer.data_buffer_file_parent_dir=/tmp/skywalking/data/buffer #缓存数据文件最大长度(单位:byte) -buffer.data_file_max_length=104857600 -#每次Flush的缓存数据的个数 -buffer.flush_number_of_cache=30 +buffer.buffer_file_max_length=104857600 #每次缓存数据写入失败,最大尝试时间 buffer.write_data_failure_retry_interval = 10000 -#数据包的最大限制 -datapackage.max_data_package=1048576 +#切换数据文件,等待时间(单位:毫秒) +persistence.switch_file_wait_time=5000 +#追加EOF标志位的线程数量 +persistence.max_append_eof_flags_thread_number=1 + +#偏移量注册文件的目录 +registerpersistence.register_file_parent_directory=/tmp/skywalking/data/offset +#偏移量注册文件名 +registerpersistence.register_file_name=offset.txt +#偏移量注册备份文件名 +registerpersistence.register_bak_file_name=offset.txt.bak +#偏移量写入文件等待周期(单位:毫秒) +registerpersistence.offset_written_file_wait_cycle=5000 + #hbase表名 hbaseconfig.table_name=sw-call-chain-new @@ -50,4 +60,4 @@ redis.edis_max_total=20 #告警检查器:异常告警检查 alarm.checker.turn_on_exception_checker=true #告警检查器:执行时间超时告警检查 -alarm.checker.turn_on_execute_time_checker=true \ No newline at end of file +alarm.checker.turn_on_execute_time_checker=true diff --git a/skywalking-server/src/test/java/test/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePluckerTest.java b/skywalking-server/src/test/java/test/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePluckerTest.java new file mode 100644 index 000000000..0f644efd4 --- /dev/null +++ b/skywalking-server/src/test/java/test/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePluckerTest.java @@ -0,0 +1,74 @@ +package test.com.ai.cloud.skywalking.reciever.peresistent; + +import com.ai.cloud.skywalking.protocol.util.IntegerAssist; +import com.ai.cloud.skywalking.reciever.peresistent.BufferBalePlucker; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.io.File; +import java.io.FileOutputStream; +import java.io.IOException; +import java.util.Arrays; + +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertNull; + +public class BufferBalePluckerTest { + + private File file; + + @Before + public void initData() throws IOException { + file = new File("/tmp", "test.file"); + if (!file.exists()) { + file.createNewFile(); + } + + FileOutputStream fileOutputStream = new FileOutputStream(file); + byte[] bytes = "HelloWorld".getBytes(); + byte[] packageByte = Arrays.copyOf(IntegerAssist.intToBytes(bytes.length), bytes.length + 4 + 4); + System.arraycopy(bytes, 0, packageByte, 4, bytes.length); + System.arraycopy(new byte[] {127, 127, 127, 127}, 0, packageByte, bytes.length + 4, 4); + + fileOutputStream.write(packageByte); + fileOutputStream.write(generatePackage("EOF".getBytes())); + fileOutputStream.close(); + } + + private static byte[] generatePackage(byte[] msg) { + byte[] dataPackage = new byte[msg.length + 8]; + // 前四位长度 + System.arraycopy(IntegerAssist.intToBytes(msg.length), 0, dataPackage, 0, 4); + // 中间正文 + System.arraycopy(msg, 0, dataPackage, 4, msg.length); + // 后四位特殊字符 + System.arraycopy(new byte[] {127, 127, 127, 127}, 0, dataPackage, msg.length + 4, 4); + + return dataPackage; + } + + @Test + public void pluck() throws Exception { + BufferBalePlucker plucker = new BufferBalePlucker(file, 0); + while (plucker.hasNextBufferBale()) { + byte[] data = plucker.pluck(); + if (data == null) + continue; + assertArrayEquals("HelloWorld".getBytes(), data); + } + } + + @Test + public void skipToNextBufferBale() throws Exception { + BufferBalePlucker plucker = new BufferBalePlucker(file, 3); + plucker.skipToNextBufferBale(); + assertNull(plucker.pluck()); + } + + @After + public void clearData() { + file.deleteOnExit(); + } + +}