package io.trygvis.test; import io.trygvis.async.SqlEffect; import io.trygvis.async.SqlEffectExecutor; import io.trygvis.queue.JdbcQueueService; import io.trygvis.queue.Queue; import io.trygvis.queue.QueueSystem; import io.trygvis.queue.Task; import io.trygvis.queue.TaskEffect; import javax.sql.DataSource; import java.sql.Connection; import java.sql.SQLException; import java.util.Date; import java.util.List; import java.util.Random; import static io.trygvis.test.DbUtil.createDataSource; import static java.lang.System.currentTimeMillis; import static java.util.Arrays.asList; import static java.util.Collections.singletonList; public class PlainJavaExample { private static final Random r = new Random(); private static String inputName = "my-input"; private static String outputName = "my-output"; private static int interval = 10; public static class Consumer { public static void main(String[] args) throws Exception { System.out.println("Starting consumer"); DataSource ds = createDataSource(); SqlEffectExecutor sqlEffectExecutor = new SqlEffectExecutor(ds); QueueSystem queueSystem = QueueSystem.initialize(sqlEffectExecutor); final JdbcQueueService queueService = queueSystem.queueService; Queue[] queues = sqlEffectExecutor.transaction(new SqlEffect() { @Override public Queue[] doInConnection(Connection c) throws SQLException { return new Queue[]{ queueService.lookupQueue(c, inputName, interval, true), queueService.lookupQueue(c, outputName, interval, true)}; } }); final Queue input = queues[0]; final Queue output = queues[1]; queueService.consumeAll(input, new TaskEffect() { public List apply(Task task) throws Exception { System.out.println("PlainJavaExample$Consumer.consumeAll: arguments = " + task.arguments); Long a = Long.valueOf(task.arguments.get(0)); Long b = Long.valueOf(task.arguments.get(1)); System.out.println("a + b = " + a + " + " + b + " = " + (a + b)); if (r.nextInt(3) == 0) { return singletonList(task.childTask(output.name, new Date(), Long.toString(a + b))); } throw new RuntimeException("Simulated exception while processing task."); } }); System.out.println("Done"); } } public static class Producer { public static void main(String[] args) throws Exception { System.out.println("Starting producer"); int chunks = 10; final int chunk = 2000; DataSource ds = createDataSource(); Connection c = ds.getConnection(); SqlEffectExecutor sqlEffectExecutor = new SqlEffectExecutor(ds); QueueSystem queueSystem = QueueSystem.initialize(sqlEffectExecutor); final JdbcQueueService queueService = queueSystem.queueService; final Queue queue = queueService.lookupQueue(c, inputName, interval, true); for (int i = 0; i < chunks; i++) { long start = currentTimeMillis(); sqlEffectExecutor.transaction(new SqlEffect.Void() { @Override public void doInConnection(Connection c) throws SQLException { for (int j = 0; j < chunk; j++) { queueService.schedule(c, queue, new Date(), asList("10", "20")); } } }); long end = currentTimeMillis(); long time = end - start; System.out.println("Scheduled " + chunk + " tasks in " + time + "ms, " + (((double) chunk * 1000)) / ((double) time) + " chunks per second"); } c.commit(); } } }