no message

This commit is contained in:
pengys5 2017-03-18 19:13:45 +08:00
parent a50171e923
commit 40fe4e5d94
17 changed files with 352 additions and 49 deletions

View File

@ -7,10 +7,4 @@ public abstract class AbstractLocalSyncWorker extends AbstractLocalWorker {
public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@Override
final public void work(Object message) throws Exception {
}
public abstract Object onWork(Object message) throws Exception;
}

View File

@ -1,18 +0,0 @@
package com.a.eye.skywalking.collector.actor;
/**
* @author pengys5
*/
public class Promise {
private boolean isTold = false;
private Object value;
protected void completed(Object value) {
this.value = value;
isTold = true;
}
public boolean isTold() {
return isTold;
}
}

View File

@ -18,13 +18,12 @@ public class TestLocalSyncWorker extends AbstractLocalSyncWorker {
}
@Override
public Object onWork(Object message) throws Exception {
public void work(Object message) throws Exception {
if (message.equals("TellLocalWorker")) {
System.out.println("hello! ");
} else {
System.out.println("unhandled");
}
return "Hello";
}
public static class Factory extends AbstractLocalSyncWorkerProvider<TestLocalSyncWorker> {

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.cluster.ClusterConfig;
import com.a.eye.skywalking.collector.cluster.ClusterConfigInitializer;
import com.a.eye.skywalking.logging.LogManager;
import com.a.eye.skywalking.logging.log4j2.Log4j2Resolver;
import com.a.eye.skywalking.collector.worker.httpserver.HttpServer;
import com.a.eye.skywalking.collector.worker.httpserver.WebServer;
import com.typesafe.config.Config;
import com.typesafe.config.ConfigFactory;
@ -32,7 +32,7 @@ public class CollectorBootStartUp {
// ActorSystem system = ActorSystem.create(ClusterConfig.Cluster.appname, config);
// WorkersCreator.INSTANCE.boot(system);
HttpServer.INSTANCE.boot();
WebServer.INSTANCE.boot();
// EsClient.boot();
}
}

View File

@ -38,7 +38,7 @@ public class ApplicationMain extends AbstractLocalSyncWorker {
}
@Override
public Object onWork(Object message) throws Exception {
public void work(Object message) throws Exception {
if (message instanceof TraceSegmentReceiver.TraceSegmentTimeSlice) {
logger.debug("begin translate TraceSegment Object to JsonObject");
TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message;
@ -50,7 +50,6 @@ public class ApplicationMain extends AbstractLocalSyncWorker {
sendToResponseCostPersistence(traceSegment);
sendToResponseSummaryPersistence(traceSegment);
}
return null;
}
public static class Factory extends AbstractLocalSyncWorkerProvider<ApplicationMain> {

View File

@ -26,7 +26,7 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker {
}
@Override
public Object onWork(Object message) throws Exception {
public void work(Object message) throws Exception {
TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message;
TraceSegment segment = traceSegment.getTraceSegment();
@ -40,7 +40,6 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker {
getSelfContext().lookup(DAGNodeRefAnalysis.Role.INSTANCE).tell(nodeRef);
}
}
return null;
}
public static class Factory extends AbstractLocalSyncWorkerProvider<ApplicationRefMain> {

View File

@ -16,8 +16,4 @@ public abstract class Controller {
protected abstract String path();
protected abstract JsonElement execute(Map<String, String> parms);
protected void tell(Role role, Object message) throws Exception {
// targetMember.beTold(message);
}
}

View File

@ -11,10 +11,10 @@ import java.util.Map;
/**
* @author pengys5
*/
public enum HttpServer {
public enum WebServer {
INSTANCE;
private Logger logger = LogManager.getFormatterLogger(HttpServer.class);
private Logger logger = LogManager.getFormatterLogger(WebServer.class);
public void boot() throws Exception {
NanoHttpServer server = new NanoHttpServer(7001);

View File

@ -0,0 +1,72 @@
package com.a.eye.skywalking.collector.worker.storage.index;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.action.admin.indices.create.CreateIndexResponse;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexResponse;
import org.elasticsearch.client.IndicesAdminClient;
import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import org.elasticsearch.index.IndexNotFoundException;
import java.io.IOException;
/**
* @author pengys5
*/
public abstract class AbstractIndex {
private Logger logger = LogManager.getFormatterLogger(AbstractIndex.class);
final public XContentBuilder createSettingBuilder() throws IOException {
XContentBuilder settingsBuilder = XContentFactory.jsonBuilder()
.startObject()
.field("index.number_of_shards", 2)
.field("index.number_of_replicas", 0)
.endObject();
return settingsBuilder;
}
public abstract XContentBuilder createMappingBuilder() throws IOException;
final public boolean createIndex() {
// settings
String settingSource = "";
// mapping
XContentBuilder mappingBuilder = null;
try {
XContentBuilder settingsBuilder = createSettingBuilder();
settingSource = settingsBuilder.string();
mappingBuilder = createMappingBuilder();
logger.info("mapping builder str: %s", mappingBuilder.string());
} catch (Exception e) {
logger.error("create %s index type of %s mapping builder error", index(), type());
}
Settings settings = Settings.builder().loadFromSource(settingSource).build();
IndicesAdminClient client = EsClient.getClient().admin().indices();
CreateIndexResponse response = client.prepareCreate(index()).setSettings(settings).addMapping(type(), mappingBuilder).get();
logger.info("create %s index with type of %s finished, isAcknowledged: %s", index(), type(), response.isAcknowledged());
return response.isAcknowledged();
}
final public boolean deleteIndex() {
IndicesAdminClient client = EsClient.getClient().admin().indices();
try {
DeleteIndexResponse response = client.prepareDelete(index()).get();
logger.info("delete %s index finished, isAcknowledged: %s", index(), response.isAcknowledged());
return response.isAcknowledged();
} catch (IndexNotFoundException e) {
logger.info("%s index not found", index());
}
return false;
}
public abstract String index();
public abstract String type();
}

View File

@ -0,0 +1,55 @@
package com.a.eye.skywalking.collector.worker.storage.index;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import java.io.IOException;
/**
* @author pengys5
*/
public class ApplicationIndexWithDagNodeType extends AbstractIndex {
private Logger logger = LogManager.getFormatterLogger(ApplicationIndexWithDagNodeType.class);
public static final String Index = "application";
public static final String Type = "dag_node";
@Override
public String index() {
return Index;
}
@Override
public String type() {
return Type;
}
@Override
public XContentBuilder createMappingBuilder() throws IOException {
XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
.startObject()
.startObject("properties")
.startObject("code")
.field("type", "string")
.field("index", "not_analyzed")
.endObject()
.startObject("layer")
.field("type", "string")
.field("index", "not_analyzed")
.endObject()
.startObject("component")
.field("type", "string")
.field("index", "not_analyzed")
.endObject()
.startObject("timeSlice")
.field("type", "long")
.field("index", "not_analyzed")
.endObject()
.endObject()
.endObject();
return mappingBuilder;
}
}

View File

@ -0,0 +1,52 @@
package com.a.eye.skywalking.collector.worker.storage.index;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import java.io.IOException;
/**
* @author pengys5
*/
public class ApplicationIndexWithNodeInstType extends AbstractIndex {
private static Logger logger = LogManager.getFormatterLogger(ApplicationIndexWithNodeInstType.class);
public static final String Index = "application";
public static final String Type = "node_instance";
@Override
public String index() {
return Index;
}
@Override
public String type() {
return Type;
}
@Override
public XContentBuilder createMappingBuilder() throws IOException {
XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
.startObject()
.startObject("properties")
.startObject("code")
.field("type", "string")
.field("fielddata", true)
.field("index", "not_analyzed")
.endObject()
.startObject("address")
.field("type", "string")
.field("index", "not_analyzed")
.endObject()
.startObject("timeSlice")
.field("type", "long_range")
.field("index", "not_analyzed")
.endObject()
.endObject()
.endObject();
return mappingBuilder;
}
}

View File

@ -0,0 +1,31 @@
package com.a.eye.skywalking.collector.worker.web.controller.tracedag;
import com.a.eye.skywalking.collector.worker.httpserver.Controller;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import fi.iki.elonen.NanoHTTPD;
import java.util.Map;
/**
* @author pengys5
*/
public class TraceDagLoadController extends Controller {
@Override
protected NanoHTTPD.Method httpMethod() {
return NanoHTTPD.Method.GET;
}
@Override
protected String path() {
return "/traceDagLoad";
}
@Override
public JsonElement execute(Map<String, String> parms) {
JsonObject jsonObject = new JsonObject();
jsonObject.addProperty("test", "aaaa");
return jsonObject;
}
}

View File

@ -0,0 +1,7 @@
package com.a.eye.skywalking.collector.worker.web.controller.tracedag;
/**
* @author pengys5
*/
public class TraceDagUpdateController {
}

View File

@ -0,0 +1,54 @@
package com.a.eye.skywalking.collector.worker.web.persistence;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.a.eye.skywalking.collector.worker.storage.index.ApplicationIndexWithDagNodeType;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.apache.logging.log4j.core.config.builder.api.FilterComponentBuilder;
import org.elasticsearch.action.search.SearchRequestBuilder;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.search.SearchType;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import org.elasticsearch.index.query.ConstantScoreQueryBuilder;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.search.aggregations.Aggregation;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.aggregations.bucket.global.GlobalAggregationBuilder;
import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import java.io.IOException;
import java.util.List;
import java.util.Map;
/**
* @author pengys5
*/
public class ApplicationPersistence {
private Logger logger = LogManager.getFormatterLogger(ApplicationPersistence.class);
public void searchDagNode(long startTimeSlice, long endTimeSlice){
SearchRequestBuilder searchRequestBuilder = EsClient.getClient().prepareSearch(ApplicationIndexWithDagNodeType.Index);
searchRequestBuilder.setTypes(ApplicationIndexWithDagNodeType.Type);
searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH);
// searchRequestBuilder.setQuery(QueryBuilders.rangeQuery("timeSlice").gte(startTimeSlice).lte(endTimeSlice));
ConstantScoreQueryBuilder constantScoreQueryBuilder = QueryBuilders.constantScoreQuery(QueryBuilders.rangeQuery("timeSlice").gte(startTimeSlice).lte(endTimeSlice));
searchRequestBuilder.setQuery(constantScoreQueryBuilder);
//// GlobalAggregationBuilder aggregationBuilder = AggregationBuilders.global("agg").subAggregation(AggregationBuilders.terms("distinct_code").field("code"));
// TermsAggregationBuilder aggregationBuilder = AggregationBuilders.terms("distinct_code").field("code");
searchRequestBuilder.addAggregation(AggregationBuilders.terms("distinct_code").field("code"));
SearchResponse response = searchRequestBuilder.execute().actionGet();
List<Aggregation> aggregationList = response.getAggregations().asList();
logger.debug("dag node list size: %s", aggregationList.size());
for (Aggregation aggregation : aggregationList) {
for (Map.Entry<String, Object> entry : aggregation.getMetaData().entrySet()) {
logger.debug("code: %s, count: %s", entry.getKey(), entry.getValue());
}
}
}
}

View File

@ -0,0 +1,7 @@
package com.a.eye.skywalking.collector.worker.web.persistence;
/**
* @author pengys5
*/
public class NodeInstancePersistence {
}

View File

@ -1,13 +1,19 @@
<?xml version="1.0" encoding="UTF-8"?>
<Configuration status="debug">
<Appenders>
<Console name="Console" target="SYSTEM_OUT">
<PatternLayout charset="UTF-8" pattern="[%d{yyyy-MM-dd HH:mm:ss:SSS}] [%p] - %l - %m%n" />
</Console>
</Appenders>
<Loggers>
<Root level="debug">
<AppenderRef ref="Console" />
</Root>
</Loggers>
<Configuration status="info">
<Appenders>
<Console name="Console" target="SYSTEM_OUT">
<PatternLayout charset="UTF-8" pattern="[%d{yyyy-MM-dd HH:mm:ss:SSS}] [%p] - %l - %m%n"/>
</Console>
</Appenders>
<Loggers>
<logger name="org.elasticsearch" level="info" additivity="false">
<AppenderRef ref="Console"/>
</logger>
<logger name="com.a.eye.skywalking.collector.worker" level="debug" additivity="false">
<AppenderRef ref="Console"/>
</logger>
<Root level="info">
<AppenderRef ref="Console"/>
</Root>
</Loggers>
</Configuration>

View File

@ -0,0 +1,50 @@
package com.a.eye.skywalking.collector.worker.web.persistence;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.a.eye.skywalking.collector.worker.storage.index.ApplicationIndexWithDagNodeType;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.rest.RestStatus;
import org.junit.Before;
import org.junit.Test;
import java.net.UnknownHostException;
import java.util.HashMap;
import java.util.Map;
/**
* @author pengys5
*/
public class ApplicationPersistenceTestCase {
@Before
public void initIndex() throws UnknownHostException {
EsClient.boot();
ApplicationIndexWithDagNodeType index = new ApplicationIndexWithDagNodeType();
index.deleteIndex();
index.createIndex();
insertData();
}
@Test
public void testSearchDagNode() {
ApplicationPersistence persistence = new ApplicationPersistence();
persistence.searchDagNode(201703101200l, 201703101259l);
}
private void insertData() {
long minute = 201703101200l;
for (int i = 0; i < 60; i++) {
Map<String, Object> json = new HashMap<String, Object>();
json.put("code", "DubboServer_MySQL");
json.put("timeSlice", minute);
json.put("component", "Dubbo");
json.put("layer", "http");
String _id = minute + "-DubboServer_MySQL";
IndexResponse response = EsClient.getClient().prepareIndex(ApplicationIndexWithDagNodeType.Index, ApplicationIndexWithDagNodeType.Type, _id).setSource(json).get();
RestStatus status = response.status();
status.getStatus();
minute++;
}
}
}