mirror of
https://gitee.com/dromara/liteFlow.git
synced 2026-05-14 20:22:07 +08:00
enhancement #I49JP1 DataBus中SlotSize的大小不支持动态扩展,无法应对高并发下的流量突增
This commit is contained in:
@@ -15,9 +15,7 @@ import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentLinkedQueue;
|
||||
import java.util.concurrent.ConcurrentSkipListSet;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicReferenceArray;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.IntStream;
|
||||
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package com.yomahub.liteflow.test.resizeSlot;
|
||||
|
||||
import cn.hutool.core.util.ReflectUtil;
|
||||
import com.yomahub.liteflow.core.FlowExecutor;
|
||||
import com.yomahub.liteflow.entity.data.DataBus;
|
||||
import com.yomahub.liteflow.entity.data.DefaultSlot;
|
||||
import com.yomahub.liteflow.entity.data.LiteflowResponse;
|
||||
import com.yomahub.liteflow.test.BaseTest;
|
||||
@@ -14,6 +16,7 @@ import org.springframework.test.context.TestPropertySource;
|
||||
import org.springframework.test.context.junit4.SpringRunner;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.*;
|
||||
@@ -39,7 +42,7 @@ public class ResizeSlotSpringbootTest extends BaseTest {
|
||||
ExecutorService pool = Executors.newCachedThreadPool();
|
||||
|
||||
List<Future<LiteflowResponse<DefaultSlot>>> futureList = new ArrayList<>();
|
||||
for (int i = 0; i < 500; i++) {
|
||||
for (int i = 0; i < 100; i++) {
|
||||
Future<LiteflowResponse<DefaultSlot>> future = pool.submit(() -> flowExecutor.execute2Resp("chain1", "arg"));
|
||||
futureList.add(future);
|
||||
}
|
||||
@@ -47,6 +50,12 @@ public class ResizeSlotSpringbootTest extends BaseTest {
|
||||
for(Future<LiteflowResponse<DefaultSlot>> future : futureList){
|
||||
Assert.assertTrue(future.get().isSuccess());
|
||||
}
|
||||
System.out.println("success");
|
||||
|
||||
//取到static的对象QUEUE
|
||||
Field field = ReflectUtil.getField(DataBus.class, "QUEUE");
|
||||
ConcurrentLinkedQueue<Integer> queue = (ConcurrentLinkedQueue<Integer>) ReflectUtil.getStaticFieldValue(field);
|
||||
|
||||
//因为初始slotSize是4,按照0.75的扩容比,要满足100个线程,应该扩容6次,6次之后应该是扩容到114
|
||||
Assert.assertEquals(queue.size(),114);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user