diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java index 7f3eff697..f74c1437d 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java @@ -12,6 +12,7 @@ import com.a.eye.skywalking.registry.impl.zookeeper.ZookeeperConfig; import com.a.eye.skywalking.storage.config.Config; import com.a.eye.skywalking.storage.config.ConfigInitializer; import com.a.eye.skywalking.storage.data.file.DataFilesManager; +import com.a.eye.skywalking.storage.data.index.operator.OperatorFactory; import com.a.eye.skywalking.storage.listener.SearchListener; import com.a.eye.skywalking.storage.listener.StorageListener; import com.a.eye.skywalking.storage.util.NetUtils; @@ -39,9 +40,10 @@ public class Main { public static void main(String[] args) { try { initializeParam(); - HealthCollector.init(SERVER_REPORTER_NAME); + OperatorFactory.initOperatorPool(); + DataFilesManager.init(); provider = ServiceProvider.newBuilder(Config.Server.PORT).addSpanStorageService(new StorageListener()) diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java index d2c520e42..6d3b4245b 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java @@ -28,7 +28,7 @@ public class Config { public static class DataIndex { - public static final int INDEX_LISTEN_PORT = 9300; + public static int INDEX_LISTEN_PORT = 9300; } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java index 0bf66fe87..93254a338 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java @@ -3,8 +3,8 @@ package com.a.eye.skywalking.storage.data; import com.a.eye.skywalking.network.grpc.TraceId; import com.a.eye.skywalking.storage.data.file.DataFileReader; import com.a.eye.skywalking.storage.data.index.*; -import com.a.eye.skywalking.storage.data.index.operator.FinderExecutor; -import com.a.eye.skywalking.storage.data.index.operator.IndexOperateExecutor; +import com.a.eye.skywalking.storage.data.index.operator.IndexOperator; +import com.a.eye.skywalking.storage.data.index.operator.OperatorFactory; import com.a.eye.skywalking.storage.data.spandata.SpanData; import java.util.ArrayList; @@ -14,8 +14,7 @@ import java.util.List; public class SpanDataFinder { public static List find(TraceId traceId) { - IndexMetaCollection indexMetaCollection = IndexOperateExecutor.execute(new FinderExecutor( - traceId.getSegmentsList().toArray(new Long[traceId.getSegmentsCount()]))); + IndexMetaCollection indexMetaCollection = fetchIndexMetaInfos(traceId); if (indexMetaCollection == null) { return new ArrayList(); @@ -45,4 +44,18 @@ public class SpanDataFinder { return result; } + + private static IndexMetaCollection fetchIndexMetaInfos(TraceId traceId) { + IndexMetaCollection indexMetaCollection = new IndexMetaCollection(); + IndexOperator indexOperator = null; + try { + indexOperator = OperatorFactory.getIndexOperatorFromPool(); + indexMetaCollection = + indexOperator.findIndex(traceId.getSegmentsList().toArray(new Long[traceId.getSegmentsCount()])); + } finally { + if (indexOperator != null) + OperatorFactory.returnIndexOperator(indexOperator); + } + return indexMetaCollection; + } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperateFailedException.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperateFailedException.java deleted file mode 100644 index aa7e0f479..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperateFailedException.java +++ /dev/null @@ -1,10 +0,0 @@ -package com.a.eye.skywalking.storage.data.exception; - -/** - * Created by xin on 2016/11/20. - */ -public class IndexOperateFailedException extends RuntimeException { - public IndexOperateFailedException(String message, Exception e) { - super(message, e); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperatorBorrowFailedException.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperatorBorrowFailedException.java new file mode 100644 index 000000000..74bfd7cd6 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/IndexOperatorBorrowFailedException.java @@ -0,0 +1,7 @@ +package com.a.eye.skywalking.storage.data.exception; + +/** + * Created by xin on 2016/11/21. + */ +public class IndexOperatorBorrowFailedException extends RuntimeException { +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java index 4a8cd6e07..89fb28be4 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java @@ -17,19 +17,19 @@ public class DataFileNameDesc { public DataFileNameDesc() { name = System.currentTimeMillis(); suffix = DATA_FILE_NAME_SUFFIX.getAndIncrement(); - fileNameStr = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SS").format(name) + "_" + suffix; + fileNameStr = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SSS").format(name) + "_" + suffix; } public DataFileNameDesc(long name, int suffix) { this.name = name; this.suffix = suffix; - fileNameStr = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SS").format(name) + "_" + suffix; + fileNameStr = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SSS").format(name) + "_" + suffix; } public DataFileNameDesc(String fileName) { int lastIndex = fileName.lastIndexOf('_'); try { - this.name = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SS").parse(fileName.substring(0, lastIndex - 1)) + this.name = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SSS").parse(fileName.substring(0, lastIndex - 1)) .getTime(); } catch (ParseException e) { } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/Executor.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/Executor.java deleted file mode 100644 index cca25918d..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/Executor.java +++ /dev/null @@ -1,5 +0,0 @@ -package com.a.eye.skywalking.storage.data.index.operator; - -interface Executor { - T execute(IndexOperator indexOperator); -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/FinderExecutor.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/FinderExecutor.java deleted file mode 100644 index ff5ff5693..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/FinderExecutor.java +++ /dev/null @@ -1,21 +0,0 @@ -package com.a.eye.skywalking.storage.data.index.operator; - - -import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; - -/** - * Created by xin on 2016/11/20. - */ -public class FinderExecutor implements Executor { - - private Long[] traceIdSegment; - - public FinderExecutor(Long[] traceIdSegment) { - this.traceIdSegment = traceIdSegment; - } - - @Override - public IndexMetaCollection execute(IndexOperator indexOperator) { - return indexOperator.findIndex(traceIdSegment); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperateExecutor.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperateExecutor.java deleted file mode 100644 index 2f4113f26..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperateExecutor.java +++ /dev/null @@ -1,36 +0,0 @@ -package com.a.eye.skywalking.storage.data.index.operator; - -import com.a.eye.skywalking.storage.config.Config; -import com.a.eye.skywalking.storage.data.exception.IndexOperateFailedException; -import com.a.eye.skywalking.storage.data.index.operator.pool.IndexOperatorPool; -import org.apache.commons.pool2.impl.GenericObjectPoolConfig; -import org.elasticsearch.client.transport.TransportClient; - -public class IndexOperateExecutor { - - private static IndexOperatorPool indexOperatorPool; - - public static T execute(Executor executor) { - TransportClient client = null; - try { - client = indexOperatorPool.borrowObject(); - return executor.execute(new IndexOperatorImpl(client)); - } catch (Exception e) { - throw new IndexOperateFailedException("Index operate failed.", e); - } finally { - indexOperatorPool.returnObject(client); - } - } - - static { - initializeIndexOperatorPool(); - } - - private static void initializeIndexOperatorPool() { - GenericObjectPoolConfig poolConfig = new GenericObjectPoolConfig(); - poolConfig.setMaxTotal(Config.IndexOperator.Finder.TOTAL); - poolConfig.setMaxIdle(Config.IndexOperator.Finder.IDEL); - poolConfig.setTestOnBorrow(true); - indexOperatorPool = new IndexOperatorPool(poolConfig); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperator.java index b767ea715..47dbaae68 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperator.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperator.java @@ -1,9 +1,96 @@ package com.a.eye.skywalking.storage.data.index.operator; +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.storage.data.file.DataFileNameDesc; import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; +import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; +import com.a.eye.skywalking.storage.data.spandata.SpanType; +import org.elasticsearch.action.bulk.BulkRequestBuilder; +import org.elasticsearch.action.bulk.BulkResponse; +import org.elasticsearch.action.search.SearchResponse; +import org.elasticsearch.client.transport.TransportClient; +import org.elasticsearch.common.xcontent.XContentBuilder; +import org.elasticsearch.index.query.BoolQueryBuilder; +import org.elasticsearch.index.query.QueryBuilders; +import org.elasticsearch.search.SearchHit; -public interface IndexOperator { - int batchUpdate(IndexMetaCollection metaInfos); +import java.io.IOException; + +import static org.elasticsearch.common.xcontent.XContentFactory.jsonBuilder; +import static org.elasticsearch.index.query.QueryBuilders.termQuery; + +public class IndexOperator { + + private static ILog logger = LogManager.getLogger(IndexOperator.class); + private final String INDEX_NAME = "skywalking"; + private final String INDEX_TYPE = "index"; + + private TransportClient client; + + public IndexOperator(TransportClient client) { + this.client = client; + } + + public int batchUpdate(IndexMetaCollection metaInfos) { + BulkRequestBuilder requestBuilder = client.prepareBulk(); + for (IndexMetaInfo indexMetaInfo : metaInfos) { + try { + requestBuilder.add(client.prepareIndex(INDEX_NAME, INDEX_TYPE).setSource(buildSource(indexMetaInfo))); + } catch (Exception e) { + logger.error("Failed to update index.", e); + HealthCollector.getCurrentHeathReading("IndexOperator") + .updateData(HeathReading.ERROR, "Failed to " + "update index."); + } + } + + BulkResponse bulkRequest = requestBuilder.get(); + if (bulkRequest.hasFailures()) { + HealthCollector.getCurrentHeathReading("IndexOperator").updateData(HeathReading.ERROR, + "Failed to " + "update index. Error message : " + bulkRequest.buildFailureMessage()); + } + + return metaInfos.size(); + } + + private XContentBuilder buildSource(IndexMetaInfo indexMetaInfo) throws IOException { + XContentBuilder xContentBuilder = jsonBuilder().startObject().field("traceid_s0", indexMetaInfo.getTraceId()[0]) + .field("traceid_s1", indexMetaInfo.getTraceId()[1]).field("traceid_s2", indexMetaInfo.getTraceId()[2]) + .field("traceid_s3", indexMetaInfo.getTraceId()[3]).field("traceid_s4", indexMetaInfo.getTraceId()[4]) + .field("traceid_s5", indexMetaInfo.getTraceId()[5]) + .field("span_type", indexMetaInfo.getSpanType().getValue()) + .field("fileName", indexMetaInfo.getFileName().getName()) + .field("fileName_suffix", indexMetaInfo.getFileName().getSuffix()) + .field("offset", indexMetaInfo.getOffset()).field("length", indexMetaInfo.getLength()).endObject(); + return xContentBuilder; + } + + public IndexMetaCollection findIndex(Long[] traceId) { + int index = 0; + BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); + for (Long traceIdSegment : traceId) { + queryBuilder.must(termQuery("traceid_s" + index++, traceIdSegment)); + } + + IndexMetaCollection collection = new IndexMetaCollection(); + SearchResponse response = + client.prepareSearch(INDEX_NAME).setTypes(INDEX_TYPE).setQuery(queryBuilder).execute().actionGet(); + for (SearchHit hit : response.getHits()) { + DataFileNameDesc desc = new DataFileNameDesc(Long.parseLong(hit.getSource().get("fileName").toString()), + Integer.parseInt(hit.getSource().get("fileName_suffix").toString())); + int length = Integer.parseInt(hit.getSource().get("length").toString()); + long offset = Long.parseLong(hit.getSource().get("offset").toString()); + SpanType spanType = SpanType.convert(Integer.parseInt(hit.getSource().get("span_type").toString())); + collection.add(new IndexMetaInfo(desc, offset, length, spanType)); + } + + return collection; + } + + public void close() { + client.close(); + } - IndexMetaCollection findIndex(Long[] traceId); } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorImpl.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorImpl.java deleted file mode 100644 index 0b3fa6c6f..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorImpl.java +++ /dev/null @@ -1,92 +0,0 @@ -package com.a.eye.skywalking.storage.data.index.operator; - -import com.a.eye.skywalking.health.report.HealthCollector; -import com.a.eye.skywalking.health.report.HeathReading; -import com.a.eye.skywalking.logging.api.ILog; -import com.a.eye.skywalking.logging.api.LogManager; -import com.a.eye.skywalking.storage.data.file.DataFileNameDesc; -import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; -import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; -import com.a.eye.skywalking.storage.data.spandata.SpanType; -import org.elasticsearch.action.bulk.BulkRequestBuilder; -import org.elasticsearch.action.bulk.BulkResponse; -import org.elasticsearch.action.search.SearchResponse; -import org.elasticsearch.client.transport.TransportClient; -import org.elasticsearch.common.xcontent.XContentBuilder; -import org.elasticsearch.index.query.BoolQueryBuilder; -import org.elasticsearch.index.query.QueryBuilder; -import org.elasticsearch.index.query.QueryBuilders; -import org.elasticsearch.search.SearchHit; - -import java.io.IOException; - -import static org.elasticsearch.common.xcontent.XContentFactory.jsonBuilder; -import static org.elasticsearch.index.query.QueryBuilders.termQuery; - -public class IndexOperatorImpl implements IndexOperator { - - private static ILog logger = LogManager.getLogger(IndexOperatorImpl.class); - - private TransportClient client; - - public IndexOperatorImpl(TransportClient client) { - this.client = client; - } - - @Override - public int batchUpdate(IndexMetaCollection metaInfos) { - BulkRequestBuilder requestBuilder = client.prepareBulk(); - for (IndexMetaInfo indexMetaInfo : metaInfos) { - try { - requestBuilder.add(client.prepareIndex("skywalking", "index").setSource(buildSource(indexMetaInfo))); - } catch (Exception e) { - logger.error("Failed to update index.", e); - HealthCollector.getCurrentHeathReading("IndexOperator") - .updateData(HeathReading.ERROR, "Failed to " + "update index."); - } - } - - BulkResponse bulkRequest = requestBuilder.get(); - if (bulkRequest.hasFailures()) { - HealthCollector.getCurrentHeathReading("IndexOperator").updateData(HeathReading.ERROR, - "Failed to " + "update index. Error message : " + bulkRequest.buildFailureMessage()); - } - - return metaInfos.size(); - } - - private XContentBuilder buildSource(IndexMetaInfo indexMetaInfo) throws IOException { - XContentBuilder xContentBuilder = jsonBuilder().startObject().field("traceid_s0", indexMetaInfo.getTraceId()[0]) - .field("traceid_s1", indexMetaInfo.getTraceId()[1]).field("traceid_s2", indexMetaInfo.getTraceId()[2]) - .field("traceid_s3", indexMetaInfo.getTraceId()[3]).field("traceid_s4", indexMetaInfo.getTraceId()[4]) - .field("traceid_s5", indexMetaInfo.getTraceId()[5]) - .field("span_type", indexMetaInfo.getSpanType().getValue()) - .field("fileName", indexMetaInfo.getFileName().getName()) - .field("fileName_suffix", indexMetaInfo.getFileName().getSuffix()) - .field("offset", indexMetaInfo.getOffset()).field("length", indexMetaInfo.getLength()).endObject(); - return xContentBuilder; - } - - @Override - public IndexMetaCollection findIndex(Long[] traceId) { - int index = 0; - BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); - for (Long traceIdSegment : traceId) { - queryBuilder.must(termQuery("traceid_s" + index++, traceIdSegment)); - } - - IndexMetaCollection collection = new IndexMetaCollection(); - SearchResponse response = client.prepareSearch("skywalking").setQuery(queryBuilder).execute().actionGet(); - for (SearchHit hit : response.getHits()) { - DataFileNameDesc desc = new DataFileNameDesc(Long.parseLong(hit.getSource().get("fileName").toString()), - Integer.parseInt(hit.getSource().get("fileName_suffix").toString())); - int length = Integer.parseInt(hit.getSource().get("length").toString()); - long offset = Long.parseLong(hit.getSource().get("offset").toString()); - SpanType spanType = SpanType.convert(Integer.parseInt(hit.getSource().get("span_type").toString())); - collection.add(new IndexMetaInfo(desc, offset, length, spanType)); - } - - return collection; - } - -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPool.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPool.java similarity index 69% rename from skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPool.java rename to skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPool.java index 8469ad507..8b6e5ee7b 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPool.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPool.java @@ -1,10 +1,10 @@ -package com.a.eye.skywalking.storage.data.index.operator.pool; +package com.a.eye.skywalking.storage.data.index.operator; import org.apache.commons.pool2.impl.GenericObjectPool; import org.apache.commons.pool2.impl.GenericObjectPoolConfig; import org.elasticsearch.client.transport.TransportClient; -public class IndexOperatorPool extends GenericObjectPool { +public class IndexOperatorPool extends GenericObjectPool { public IndexOperatorPool(GenericObjectPoolConfig poolConfig) { super(new IndexOperatorPooledObjectFactory(), poolConfig); } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPooledObjectFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPooledObjectFactory.java similarity index 55% rename from skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPooledObjectFactory.java rename to skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPooledObjectFactory.java index 85fa34318..26e5eb337 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/pool/IndexOperatorPooledObjectFactory.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/IndexOperatorPooledObjectFactory.java @@ -1,4 +1,4 @@ -package com.a.eye.skywalking.storage.data.index.operator.pool; +package com.a.eye.skywalking.storage.data.index.operator; import com.a.eye.skywalking.storage.config.Config; import org.apache.commons.pool2.BasePooledObjectFactory; @@ -11,21 +11,20 @@ import org.elasticsearch.transport.client.PreBuiltTransportClient; import java.net.InetAddress; -public class IndexOperatorPooledObjectFactory extends BasePooledObjectFactory { +public class IndexOperatorPooledObjectFactory extends BasePooledObjectFactory { @Override - public TransportClient create() throws Exception { - return new PreBuiltTransportClient(Settings.EMPTY).addTransportAddress( - new InetSocketTransportAddress(InetAddress.getLocalHost(), Config.DataIndex.INDEX_LISTEN_PORT)); + public IndexOperator create() throws Exception { + return OperatorFactory.createIndexOperator(); } @Override - public PooledObject wrap(TransportClient client) { + public PooledObject wrap(IndexOperator client) { return new DefaultPooledObject<>(client); } @Override - public void destroyObject(PooledObject p) throws Exception { - super.destroyObject(p); + public void destroyObject(PooledObject p) throws Exception { + p.getObject().close(); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/OperatorFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/OperatorFactory.java index a229a0d6d..dc2496a0a 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/OperatorFactory.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/OperatorFactory.java @@ -1,7 +1,9 @@ package com.a.eye.skywalking.storage.data.index.operator; import com.a.eye.skywalking.storage.config.Config; +import com.a.eye.skywalking.storage.data.exception.IndexOperatorBorrowFailedException; import com.a.eye.skywalking.storage.data.exception.IndexOperatorInitializeFailedException; +import org.apache.commons.pool2.impl.GenericObjectPoolConfig; import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.transport.InetSocketTransportAddress; import org.elasticsearch.transport.client.PreBuiltTransportClient; @@ -10,13 +12,34 @@ import java.net.InetAddress; public class OperatorFactory { + private static IndexOperatorPool pool; + public static IndexOperator createIndexOperator() { try { - return new IndexOperatorImpl(new PreBuiltTransportClient(Settings.EMPTY).addTransportAddress( + return new IndexOperator(new PreBuiltTransportClient(Settings.EMPTY).addTransportAddress( new InetSocketTransportAddress(InetAddress.getLocalHost(), Config.DataIndex.INDEX_LISTEN_PORT))); } catch (Exception e) { throw new IndexOperatorInitializeFailedException("Failed to initialize operator.", e); } } + public static IndexOperator getIndexOperatorFromPool() { + try { + return pool.borrowObject(); + } catch (Exception e) { + throw new IndexOperatorBorrowFailedException(); + } + } + + public static void returnIndexOperator(IndexOperator indexOperator) { + pool.returnObject(indexOperator); + } + + public static void initOperatorPool() { + GenericObjectPoolConfig config = new GenericObjectPoolConfig(); + config.setMaxTotal(Config.IndexOperator.Finder.TOTAL); + config.setMaxIdle(Config.IndexOperator.Finder.IDEL); + pool = new IndexOperatorPool(config); + } + } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/UpdateExecutor.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/UpdateExecutor.java deleted file mode 100644 index e779267e1..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/operator/UpdateExecutor.java +++ /dev/null @@ -1,17 +0,0 @@ -package com.a.eye.skywalking.storage.data.index.operator; - -import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; - -public class UpdateExecutor implements Executor { - - private IndexMetaCollection metaCollection; - - public UpdateExecutor(IndexMetaCollection metaCollection) { - this.metaCollection = metaCollection; - } - - @Override - public Integer execute(IndexOperator indexOperator) { - return indexOperator.batchUpdate(metaCollection); - } -}