diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java
index 207075aba..4adff6da6 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java
@@ -11,18 +11,44 @@ import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
/**
+ * The AbstractClusterWorker implementations represent workers,
+ * which receive remote messages.
+ *
+ * Usually, the implementations are doing persistent, or aggregate works.
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public abstract class AbstractClusterWorker extends AbstractWorker {
+ /**
+ * Construct an AbstractClusterWorker with the worker role and context.
+ *
+ * @param role If multi-workers are for load balance, they should be more likely called worker instance.
+ * Meaning, each worker have multi instances.
+ * @param clusterContext See {@link ClusterWorkerContext}
+ * @param selfContext See {@link LocalWorkerContext}
+ */
protected AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
+ /**
+ * This method use for message producer to call for send message.
+ *
+ * @param message The persistence data or metric data.
+ * @throws Exception The Exception happen in {@link #onWork(Object)}
+ */
final public void allocateJob(Object message) throws Exception {
onWork(message);
}
+ /**
+ * This method use for message receiver to analyse message.
+ *
+ * @param message Cast the message object to a expect subclass.
+ * @throws Exception Don't handle the exception, throw it.
+ */
protected abstract void onWork(Object message) throws Exception;
static class WorkerWithAkka extends UntypedActor {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java
index 421d44112..6b7bb5ef5 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java
@@ -4,17 +4,36 @@ import akka.actor.ActorRef;
import akka.actor.Props;
/**
+ * The AbstractClusterWorkerProvider implementations represent providers,
+ * which create instance of cluster workers whose implemented {@link AbstractClusterWorker}.
+ *
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider {
+ /**
+ * Create how many worker instance of {@link AbstractClusterWorker} in one jvm.
+ *
+ * @return The worker instance number.
+ */
public abstract int workerNum();
+ /**
+ * Create the worker instance into akka system, the akka system will control the cluster worker life cycle.
+ *
+ * @param localContext Not used, will be null.
+ * @return The created worker reference. See {@link ClusterWorkerRef}
+ * @throws IllegalArgumentException Not used.
+ * @throws ProviderNotFoundException This worker instance attempted to find a provider which use to create another worker
+ * instance, when the worker provider not find then Throw this Exception.
+ */
@Override
final public WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFoundException {
int num = ClusterWorkerRefCounter.INSTANCE.incrementAndGet(role());
- T clusterWorker = (T) workerInstance(getClusterContext());
+ T clusterWorker = workerInstance(getClusterContext());
clusterWorker.preStart();
ActorRef actorRef = getClusterContext().getAkkaSystem().actorOf(Props.create(AbstractClusterWorker.WorkerWithAkka.class, clusterWorker), role().roleName() + "_" + num);
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java
index eb919e57e..915128af9 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java
@@ -6,34 +6,72 @@ import com.lmax.disruptor.EventHandler;
import com.lmax.disruptor.RingBuffer;
/**
+ * The AbstractLocalAsyncWorker implementations represent workers,
+ * which receive local asynchronous message.
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
+ /**
+ * Construct an AbstractLocalAsyncWorker with the worker role and context.
+ *
+ * @param role The responsibility of worker in cluster, more than one workers can have
+ * same responsibility which use to provide load balancing ability.
+ * @param clusterContext See {@link ClusterWorkerContext}
+ * @param selfContext See {@link LocalWorkerContext}
+ */
public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
+ /**
+ * The asynchronous worker always use to persistence data into db, this is the end of the streaming,
+ * so usually no need to create the next worker instance at the time of this worker instance create.
+ *
+ * @throws ProviderNotFoundException When worker provider not found, it will be throw this exception.
+ */
@Override
public void preStart() throws ProviderNotFoundException {
}
- final public void allocateJob(Object request) throws Exception {
- onWork(request);
+ /**
+ * Receive message
+ *
+ * @param message The persistence data or metric data.
+ * @throws Exception The Exception happen in {@link #onWork(Object)}
+ */
+ final public void allocateJob(Object message) throws Exception {
+ onWork(message);
}
- protected abstract void onWork(Object request) throws Exception;
+ /**
+ * The data process logic in this method.
+ *
+ * @param message Cast the message object to a expect subclass.
+ * @throws Exception Don't handle the exception, throw it.
+ */
+ protected abstract void onWork(Object message) throws Exception;
static class WorkerWithDisruptor implements EventHandler {
private RingBuffer ringBuffer;
private AbstractLocalAsyncWorker asyncWorker;
- public WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) {
+ WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) {
this.ringBuffer = ringBuffer;
this.asyncWorker = asyncWorker;
}
+ /**
+ * Receive the message from disruptor, when message in disruptor is empty, then send the cached data
+ * to the next workers.
+ *
+ * @param event published to the {@link RingBuffer}
+ * @param sequence of the event being processed
+ * @param endOfBatch flag to indicate if this is the last event in a batch from the {@link RingBuffer}
+ */
public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) {
try {
Object message = event.getMessage();
@@ -48,6 +86,12 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
}
}
+ /**
+ * Push the message into disruptor ring buffer.
+ *
+ * @param message of the data to process.
+ * @throws Exception not used.
+ */
public void tell(Object message) throws Exception {
long sequence = ringBuffer.next();
try {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java
index 0a5f64a14..df1449b5d 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java
@@ -1,7 +1,13 @@
package com.a.eye.skywalking.collector.actor.selector;
/**
+ * The AbstractHashMessage implementations represent aggregate message,
+ * which use to aggregate metric.
+ * Make the message aggregator's worker selector use of {@link HashCodeSelector}.
+ *
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public abstract class AbstractHashMessage {
private int hashCode;
@@ -10,7 +16,7 @@ public abstract class AbstractHashMessage {
this.hashCode = key.hashCode();
}
- protected int getHashCode() {
+ int getHashCode() {
return hashCode;
}
}
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java
index a3d46df68..62a21026d 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java
@@ -1,14 +1,27 @@
package com.a.eye.skywalking.collector.actor.selector;
+import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.WorkerRef;
import java.util.List;
/**
+ * The HashCodeSelector is a simple implementation of {@link WorkerSelector}.
+ * It choose {@link WorkerRef} by message {@link AbstractHashMessage} key's hashcode, so it can use to send the same hashcode
+ * message to same {@link WorkerRef}. Usually, use to database operate which avoid dirty data.
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public class HashCodeSelector implements WorkerSelector {
+ /**
+ * Use message hashcode to select {@link WorkerRef}.
+ *
+ * @param members given {@link WorkerRef} list, which size is greater than 0;
+ * @param message the {@link AbstractWorker} is going to send.
+ * @return the selected {@link WorkerRef}
+ */
@Override
public WorkerRef select(List members, Object message) {
if (message instanceof AbstractHashMessage) {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java
index ec0f89822..256519832 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java
@@ -1,16 +1,28 @@
package com.a.eye.skywalking.collector.actor.selector;
+import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.WorkerRef;
import java.util.List;
/**
+ * The RollingSelector is a simple implementation of {@link WorkerSelector}.
+ * It choose {@link WorkerRef} nearly random, by round-robin.
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public class RollingSelector implements WorkerSelector {
private int index = 0;
+ /**
+ * Use round-robin to select {@link WorkerRef}.
+ *
+ * @param members given {@link WorkerRef} list, which size is greater than 0;
+ * @param message message the {@link AbstractWorker} is going to send.
+ * @return the selected {@link WorkerRef}
+ */
@Override
public WorkerRef select(List members, Object message) {
int size = members.size();
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java
index c7a607f47..7ea1687a7 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java
@@ -1,12 +1,27 @@
package com.a.eye.skywalking.collector.actor.selector;
+import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.WorkerRef;
import java.util.List;
/**
+ * The WorkerSelector should be implemented by any class whose instances
+ * are intended to provide select a {@link WorkerRef} from a {@link WorkerRef} list.
+ *
+ * Actually, the WorkerRef is designed to provide a routing ability in the collector cluster
+ *
* @author pengys5
+ * @since v3.0-2017
*/
public interface WorkerSelector {
+
+ /**
+ * select a {@link WorkerRef} from a {@link WorkerRef} list.
+ *
+ * @param members given {@link WorkerRef} list, which size is greater than 0;
+ * @param message the {@link AbstractWorker} is going to send.
+ * @return the selected {@link WorkerRef}
+ */
T select(List members, Object message);
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/config/EsConfig.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/config/EsConfig.java
index 426307850..ae6ffb4af 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/config/EsConfig.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/config/EsConfig.java
@@ -16,13 +16,22 @@ public class EsConfig {
}
public static class Index {
- public static class Shards {
- public static String number;
+
+ public static class Initialize {
+ public static IndexInitMode mode;
}
- public static class Replicas{
- public static String number;
+ public static class Shards {
+ public static String number = "";
+ }
+
+ public static class Replicas {
+ public static String number = "";
}
}
}
+
+ public enum IndexInitMode {
+ auto, forced, manual
+ }
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndex.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndex.java
index 1aef5a01d..fa16d5760 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndex.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndex.java
@@ -5,6 +5,7 @@ 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.action.admin.indices.exists.indices.IndicesExistsResponse;
import org.elasticsearch.action.admin.indices.mapping.put.PutMappingResponse;
import org.elasticsearch.client.IndicesAdminClient;
import org.elasticsearch.common.settings.Settings;
@@ -86,5 +87,11 @@ public abstract class AbstractIndex {
return false;
}
+ final boolean isExists() {
+ IndicesAdminClient client = EsClient.INSTANCE.getClient().admin().indices();
+ IndicesExistsResponse response = client.prepareExists(index()).get();
+ return response.isExists();
+ }
+
public abstract String index();
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/EsClient.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/EsClient.java
index b4985ce24..969721de9 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/EsClient.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/EsClient.java
@@ -13,6 +13,7 @@ import org.elasticsearch.transport.client.PreBuiltTransportClient;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
/**
@@ -46,9 +47,9 @@ public enum EsClient {
public void indexRefresh(String... indexName) {
RefreshResponse response = client.admin().indices().refresh(new RefreshRequest(indexName)).actionGet();
if (response.getShardFailures().length == response.getTotalShards()) {
- logger.error("All elasticsearch shard index refresh failure, reason: %s", response.getShardFailures());
+ logger.error("All elasticsearch shard index refresh failure, reason: %s", Arrays.toString(response.getShardFailures()));
} else if (response.getShardFailures().length > 0) {
- logger.error("In parts of elasticsearch shard index refresh failure, reason: %s", response.getShardFailures());
+ logger.error("In parts of elasticsearch shard index refresh failure, reason: %s", Arrays.toString(response.getShardFailures()));
}
logger.info("elasticsearch index refresh success");
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/IndexCreator.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/IndexCreator.java
index 1c2073cdd..973ea1ee0 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/IndexCreator.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/IndexCreator.java
@@ -1,5 +1,6 @@
package com.a.eye.skywalking.collector.worker.storage;
+import com.a.eye.skywalking.collector.worker.config.EsConfig;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@@ -16,10 +17,19 @@ public enum IndexCreator {
private Logger logger = LogManager.getFormatterLogger(IndexCreator.class);
public void create() {
- Set indexSet = loadIndex();
- for (AbstractIndex index : indexSet) {
- index.deleteIndex();
- index.createIndex();
+ if (!EsConfig.IndexInitMode.manual.equals(EsConfig.Es.Index.Initialize.mode)) {
+ Set indexSet = loadIndex();
+ for (AbstractIndex index : indexSet) {
+ boolean isExists = index.isExists();
+ if (isExists) {
+ if (EsConfig.IndexInitMode.forced.equals(EsConfig.Es.Index.Initialize.mode)) {
+ index.deleteIndex();
+ index.createIndex();
+ }
+ } else {
+ index.createIndex();
+ }
+ }
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config b/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config
index 0c35dd850..dfe77c631 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config
+++ b/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config
@@ -1,20 +1,52 @@
+# The remote server should connect to, hostname can be either hostname or IP address.
+# Suggestion: set the real ip address.
cluster.current.hostname=127.0.0.1
cluster.current.port=11800
+
+# The roles of this member. List of strings, e.g. roles = A, B
+# In the future, the roles are part of the membership information and can be used by
+# routers or other services to distribute work to certain member types,
+# e.g. front-end and back-end nodes.
+# In this version, all members has same roles, each of them will listen others status,
+# because of network trouble or member jvm crash or every reason led to not reachable,
+# the routers will stop to sending the message to the untouchable member.
cluster.current.roles=WorkersListener
+
+# Initial contact points of the cluster, e.g. seed_nodes = 127.0.0.1:11800, 127.0.0.1:11801.
+# The nodes to join automatically at startup.
+# When setting akka configuration, it will be change.
+# like: ["akka.tcp://system@127.0.0.1:11800", "akka.tcp://system@127.0.0.1:11801"].
+# This is akka configuration, see: http://doc.akka.io/docs/akka/2.4/general/configuration.html
cluster.seed_nodes=127.0.0.1:11800
+# elasticsearch configuration, config/elasticsearch.yml, see cluster.name
es.cluster.name=CollectorDBCluster
-es.cluster.nodes=127.0.0.1:9300
es.cluster.transport.sniffer=true
+# The elasticsearch nodes of cluster, comma separated, e.g. nodes=ip:port, ip:port
+es.cluster.nodes=127.0.0.1:9300
+
+# Automatic create elasticsearch index
+# Options: auto, forced, manual
+# auto: just create new index when index not created.
+# forced: delete the index then create
+es.index.initialize.mode=auto
es.index.shards.number=2
es.index.replicas.number=0
+# You can configure a host either as a host name or IP address to identify a specific network
+# interface on which to listen.
+# Be used for web ui get the view data or agent post the trace segment.
http.hostname=127.0.0.1
+# The TCP/IP port on which the connector listens for connections.
http.port=12800
+# The contextPath is a URL prefix that identifies which context a HTTP request is destined for.
http.contextPath=/
+# The analysis worker max cache size, when worker data size reach the size,
+# then worker will send all cached data to the next worker and clear the cache.
cache.analysis.size=1024
+# The persistence worker max cache size, same of "cache.analysis.size" ability.
cache.persistence.size=1024
WorkerNum.Node.NodeCompAgg.Value=10
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndexTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndexTestCase.java
index f53a8abf6..47f20a1c6 100644
--- a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndexTestCase.java
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/AbstractIndexTestCase.java
@@ -22,7 +22,7 @@ public class AbstractIndexTestCase {
@Test
public void testCreateSettingBuilder() throws IOException {
IndexTest indexTest = new IndexTest();
- Assert.assertEquals("{\"index.number_of_shards\":null,\"index.number_of_replicas\":null}", indexTest.createSettingBuilder().string());
+ Assert.assertEquals("{\"index.number_of_shards\":\"\",\"index.number_of_replicas\":\"\"}", indexTest.createSettingBuilder().string());
}
class IndexTest extends AbstractIndex {
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/IndexCreatorTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/IndexCreatorTestCase.java
index f4c64e100..52eaf2c1c 100644
--- a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/IndexCreatorTestCase.java
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/storage/IndexCreatorTestCase.java
@@ -1,8 +1,10 @@
package com.a.eye.skywalking.collector.worker.storage;
+import com.a.eye.skywalking.collector.worker.config.EsConfig;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import org.junit.Assert;
+import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mockito;
@@ -25,6 +27,21 @@ import static org.powermock.api.mockito.PowerMockito.*;
@PowerMockIgnore({"javax.management.*"})
public class IndexCreatorTestCase {
+ private IndexCreator indexCreator;
+ private TestIndex testIndex;
+
+ @Before
+ public void init() throws Exception {
+ testIndex = mock(TestIndex.class);
+
+ indexCreator = mock(IndexCreator.class);
+ doCallRealMethod().when(indexCreator).create();
+
+ Set indexSet = new HashSet<>();
+ indexSet.add(testIndex);
+ when(indexCreator, "loadIndex").thenReturn(indexSet);
+ }
+
@Test
public void testLoadIndex() throws Exception {
IndexCreator indexCreator = spy(IndexCreator.INSTANCE);
@@ -47,22 +64,49 @@ public class IndexCreatorTestCase {
}
@Test
- public void testCreate() throws Exception {
- TestIndex testIndex = mock(TestIndex.class);
-
- IndexCreator indexCreator = mock(IndexCreator.class);
- doCallRealMethod().when(indexCreator).create();
-
- Set indexSet = new HashSet<>();
- indexSet.add(testIndex);
-
- when(indexCreator, "loadIndex").thenReturn(indexSet);
+ public void testCreateOptionManual() throws Exception {
+ EsConfig.Es.Index.Initialize.mode = EsConfig.IndexInitMode.manual;
+ indexCreator.create();
+ Mockito.verify(testIndex, Mockito.never()).createIndex();
+ Mockito.verify(testIndex, Mockito.never()).deleteIndex();
+ }
+ @Test
+ public void testCreateOptionForcedIndexIsExists() throws Exception {
+ EsConfig.Es.Index.Initialize.mode = EsConfig.IndexInitMode.forced;
+ when(testIndex.isExists()).thenReturn(true);
indexCreator.create();
Mockito.verify(testIndex).createIndex();
Mockito.verify(testIndex).deleteIndex();
}
+ @Test
+ public void testCreateOptionForcedIndexNotExists() throws Exception {
+ EsConfig.Es.Index.Initialize.mode = EsConfig.IndexInitMode.forced;
+ when(testIndex.isExists()).thenReturn(false);
+ indexCreator.create();
+ Mockito.verify(testIndex).createIndex();
+ Mockito.verify(testIndex, Mockito.never()).deleteIndex();
+ }
+
+ @Test
+ public void testCreateOptionAutoIndexNotExists() throws Exception {
+ EsConfig.Es.Index.Initialize.mode = EsConfig.IndexInitMode.auto;
+ when(testIndex.isExists()).thenReturn(false);
+ indexCreator.create();
+ Mockito.verify(testIndex).createIndex();
+ Mockito.verify(testIndex, Mockito.never()).deleteIndex();
+ }
+
+ @Test
+ public void testCreateOptionAutoIndexExists() throws Exception {
+ EsConfig.Es.Index.Initialize.mode = EsConfig.IndexInitMode.auto;
+ when(testIndex.isExists()).thenReturn(true);
+ indexCreator.create();
+ Mockito.verify(testIndex, Mockito.never()).createIndex();
+ Mockito.verify(testIndex, Mockito.never()).deleteIndex();
+ }
+
class TestIndex extends AbstractIndex {
@Override