feat: add normal & manual bus

This commit is contained in:
tiennm99 committed 2025-02-25 21:09:56 +07:00
1 parent b512add919
commit 53ef817551
14 files changed
+198 -10

No files matched your search

+11
View File
@@ -1,2 +1,13 @@
# taskbus
A task distribution, which assign tasks to specific executors
## How to use
Just view tests in [`test`](src/test/java) folder for example.
## Caution
This project is just demo a simple task distribution.
In production there will be more complicated cases that you need to handle, like executor type,
monitor, etc.
+24 -3
View File
@@ -1,19 +1,40 @@
plugins {
id("java")
java
idea
}
group = "com.miti99"
version = "1.0-SNAPSHOT"
java {
toolchain {
languageVersion.set(JavaLanguageVersion.of(21))
}
}
idea {
module {
isDownloadJavadoc = true
isDownloadSources = true
}
}
repositories {
mavenCentral()
}
dependencies {
testImplementation(platform("org.junit:junit-bom:5.10.0"))
annotationProcessor("org.projectlombok:lombok:1.18.36")
compileOnly("org.projectlombok:lombok:1.18.36")
implementation("org.slf4j:slf4j-simple:2.0.16")
testAnnotationProcessor("org.projectlombok:lombok:1.18.36")
testCompileOnly("org.projectlombok:lombok:1.18.36")
testImplementation(platform("org.junit:junit-bom:5.11.4"))
testImplementation("org.junit.jupiter:junit-jupiter")
}
tasks.test {
useJUnitPlatform()
}
}
@@ -1,7 +0,0 @@
package com.miti99;
public class Main {
public static void main(String[] args) {
System.out.println("Hello, World!");
}
}
@@ -0,0 +1,7 @@
package com.miti99.taskbus.bus;
import com.miti99.taskbus.task.Task;
public interface Bus {
void submit(Task task);
}
@@ -0,0 +1,29 @@
package com.miti99.taskbus.bus.impl;
import com.miti99.taskbus.bus.Bus;
import com.miti99.taskbus.task.Task;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class ManualBus implements Bus {
private final ExecutorService[] executors;
private final int poolSize;
public ManualBus(int poolSize) {
this.poolSize = poolSize;
executors = new ExecutorService[poolSize];
for (int i = 0; i < poolSize; i++) {
executors[i] = Executors.newSingleThreadExecutor();
}
}
public ManualBus() {
this(Runtime.getRuntime().availableProcessors());
}
@Override
public void submit(Task task) {
int executorId = task.hash() % poolSize;
executors[executorId].submit(task::execute);
}
}
@@ -0,0 +1,16 @@
package com.miti99.taskbus.bus.impl;
import com.miti99.taskbus.bus.Bus;
import com.miti99.taskbus.task.Task;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class NormalBus implements Bus {
private final ExecutorService executors =
Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
@Override
public void submit(Task task) {
executors.submit(task::execute);
}
}
@@ -0,0 +1,7 @@
package com.miti99.taskbus.task;
public interface Task {
int hash();
void execute();
}
@@ -0,0 +1,12 @@
package com.miti99.taskbus.task.event;
import com.miti99.taskbus.task.Task;
public interface Event extends Task {
int userId();
@Override
default int hash() {
return userId();
}
}
@@ -0,0 +1,12 @@
package com.miti99.taskbus.task.event;
import lombok.extern.slf4j.Slf4j;
@Slf4j
public record UserLoginEvent(int userId) implements Event {
@Override
public void execute() {
log.info("User {} logged in", userId);
}
}
@@ -0,0 +1,12 @@
package com.miti99.taskbus.task.request;
import com.miti99.taskbus.task.Task;
public interface Request extends Task {
int userId();
@Override
default int hash() {
return userId();
}
}
@@ -0,0 +1,12 @@
package com.miti99.taskbus.task.request;
import lombok.extern.slf4j.Slf4j;
@Slf4j
public record SendMessageRequest(int userId, String message) implements Request {
@Override
public void execute() {
log.info("User {} sent message: {}", userId, message);
}
}
@@ -0,0 +1,24 @@
package com.miti99.taskbus.bus.impl;
import com.miti99.taskbus.bus.Bus;
import com.miti99.taskbus.task.event.UserLoginEvent;
import com.miti99.taskbus.task.request.SendMessageRequest;
import lombok.RequiredArgsConstructor;
@RequiredArgsConstructor
public abstract class BusTest {
private final Bus bus;
void submitTasks() {
bus.submit(new UserLoginEvent(1));
bus.submit(new SendMessageRequest(1, "msg1"));
bus.submit(new SendMessageRequest(1, "msg2"));
bus.submit(new UserLoginEvent(2));
bus.submit(new SendMessageRequest(2, "msg1"));
bus.submit(new SendMessageRequest(2, "msg2"));
bus.submit(new SendMessageRequest(1, "msg3"));
bus.submit(new SendMessageRequest(1, "msg4"));
bus.submit(new SendMessageRequest(2, "msg3"));
bus.submit(new SendMessageRequest(2, "msg4"));
}
}
@@ -0,0 +1,15 @@
package com.miti99.taskbus.bus.impl;
import org.junit.jupiter.api.Test;
class ManualBusTest extends BusTest {
public ManualBusTest() {
super(new ManualBus());
}
@Test
void runTest() {
submitTasks();
}
}
@@ -0,0 +1,17 @@
package com.miti99.taskbus.bus.impl;
import static org.junit.jupiter.api.Assertions.*;
import org.junit.jupiter.api.Test;
class NormalBusTest extends BusTest {
public NormalBusTest() {
super(new NormalBus());
}
@Test
void runTest() {
submitTasks();
}
}