phantomthief / buffer-trigger Goto Github PK
View Code? Open in Web Editor NEWLicense: Artistic License 2.0
License: Artistic License 2.0
可以稳定复现的场景是这样的:
public class TestTrigger {
private AtomicLong enqueueCount = new AtomicLong();
private AtomicLong consumeCount = new AtomicLong();
private AtomicLong rejectCount = new AtomicLong();
private BufferTrigger<String> buffer = BufferTrigger.<String, Queue<String>>simple()
.name("test-trigger")
.setContainer(ConcurrentLinkedQueue::new, Queue::add)
.maxBufferCount(1000)
.interval(1, TimeUnit.SECONDS)
.consumer(this::doBatchReload)
.rejectHandler(this::onTaskRejected)
.build();
private void doBatchReload(Iterable<String> values) {
consumeCount.addAndGet(Iterables.size(values));
Uninterruptibles.sleepUninterruptibly(Duration.ofMillis(1000));
}
private void onTaskRejected(String value) {
rejectCount.addAndGet(1);
}
private void test() throws InterruptedException {
ExecutorService executor = Executors.newFixedThreadPool(100);
for (int i = 0; i < 1000000; i++) {
executor.submit(() -> {
enqueueCount.getAndAdd(1);
buffer.enqueue("test");
});
if (i % 353 == 0) {
Uninterruptibles.sleepUninterruptibly(Duration.ofMillis(50));
}
}
executor.shutdown();
boolean finished = executor.awaitTermination(30, TimeUnit.SECONDS);
System.out.println(finished);
buffer.manuallyDoTrigger();
System.out.printf("enqueued: %d\n", enqueueCount.get());
System.out.printf("handled: %d + %d = %d\n", consumeCount.get(), rejectCount.get(), consumeCount.get() + rejectCount.get());
}
public static void main(String[] args) throws InterruptedException {
TestTrigger test = new TestTrigger();
test.test();
}
}
结果是:
true
enqueued: 1000000
handled: 150023 + 849973 = 999996
因为每台机器都会enqueue() 这也意味着每台机器都会触发Trigger的逻辑, 只能用在MQ的消费逻辑的enqueue() 做别的业务批量处理逻辑吗?
version:
lastest
JVM version (java -version):
1.8
Description of the problem including expected versus actual behavior:
I constructed BatchConsumeBlockingQueueTrigger with the following arguments:
batchSize:1000 lingerMs:5000
The enqueue qps is 200/s and it spent 50ms to consume the bulk data despite the number of the bulk data.
expected behavior
the consumer function invoke every 5 seconds with about 1000 data.
actual behavior
the consumer function invoked more than 5 times in 1 seconds and the number of bulk data less than 200.
version:
lastest
JVM version (java -version):
1.8
Description of the problem including expected versus actual behavior:
expected behavior
consumer function is called when the batchSize is reached or the linger time passed
actual behavior
consumer function is only callen when linger time passed
code to reproduce
public static void main(String[] args) {
ScheduledExecutorService threadPoolExecutorService = new ScheduledThreadPoolExecutor(1);
ScheduledExecutorService scheduledExecutorService = new ScheduledExecutorService() {
@NotNull
@Override
public ScheduledFuture<?> schedule(@NotNull Runnable command, long delay, @NotNull TimeUnit unit) {
return threadPoolExecutorService.schedule(command, delay, unit);
}
@NotNull
@Override
public <V> ScheduledFuture<V> schedule(@NotNull Callable<V> callable, long delay, @NotNull TimeUnit unit) {
return threadPoolExecutorService.schedule(callable, delay, unit);
}
@NotNull
@Override
public ScheduledFuture<?> scheduleAtFixedRate(@NotNull Runnable command, long initialDelay, long period, @NotNull TimeUnit unit) {
return threadPoolExecutorService.scheduleAtFixedRate(command, initialDelay, period, unit);
}
@NotNull
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(@NotNull Runnable command, long initialDelay, long delay, @NotNull TimeUnit unit) {
return null;
}
@Override
public void shutdown() {
}
@NotNull
@Override
public List<Runnable> shutdownNow() {
return null;
}
@Override
public boolean isShutdown() {
return false;
}
@Override
public boolean isTerminated() {
return false;
}
@Override
public boolean awaitTermination(long timeout, @NotNull TimeUnit unit) throws InterruptedException {
return false;
}
@NotNull
@Override
public <T> Future<T> submit(@NotNull Callable<T> task) {
return null;
}
@NotNull
@Override
public <T> Future<T> submit(@NotNull Runnable task, T result) {
return null;
}
@NotNull
@Override
public Future<?> submit(@NotNull Runnable task) {
return null;
}
@NotNull
@Override
public <T> List<Future<T>> invokeAll(@NotNull Collection<? extends Callable<T>> tasks) throws InterruptedException {
return null;
}
@NotNull
@Override
public <T> List<Future<T>> invokeAll(@NotNull Collection<? extends Callable<T>> tasks, long timeout, @NotNull TimeUnit unit) throws InterruptedException {
return null;
}
@NotNull
@Override
public <T> T invokeAny(@NotNull Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException {
return null;
}
@Override
public <T> T invokeAny(@NotNull Collection<? extends Callable<T>> tasks, long timeout, @NotNull TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
return null;
}
@Override
public void execute(@NotNull Runnable command) {
command.run();
}
};
BufferTrigger<String> bufferTrigger = BufferTrigger.<String>batchBlocking()
.batchSize(20).linger(3000, TimeUnit.SECONDS)
.setConsumerEx(t -> {
System.out.println(Arrays.toString(t.toArray()));
}).setScheduleExecutorService(scheduledExecutorService).build();
for (int j = 0; j < 100; j++) {
for (int i = 0; i < 10; i++) {
bufferTrigger.enqueue(String.valueOf(j * 100 + i));
}
System.out.println("pending:" + bufferTrigger.getPendingChanges());
}
Uninterruptibles.sleepUninterruptibly(3, TimeUnit.HOURS);
}
Related code
change the running.set(true); before scheduledExecutorService.execute?
A declarative, efficient, and flexible JavaScript library for building user interfaces.
🖖 Vue.js is a progressive, incrementally-adoptable JavaScript framework for building UI on the web.
TypeScript is a superset of JavaScript that compiles to clean JavaScript output.
An Open Source Machine Learning Framework for Everyone
The Web framework for perfectionists with deadlines.
A PHP framework for web artisans
Bring data to life with SVG, Canvas and HTML. 📊📈🎉
JavaScript (JS) is a lightweight interpreted programming language with first-class functions.
Some thing interesting about web. New door for the world.
A server is a program made to process requests and deliver data to clients.
Machine learning is a way of modeling and interpreting data that allows a piece of software to respond intelligently.
Some thing interesting about visualization, use data art
Some thing interesting about game, make everyone happy.
We are working to build community through open source technology. NB: members must have two-factor auth.
Open source projects and samples from Microsoft.
Google ❤️ Open Source for everyone.
Alibaba Open Source for everyone
Data-Driven Documents codes.
China tencent open source team.