Notification Service
Delayed sends, idempotent enqueues, priority lanes, periodic digests, and batch fan-out for a notification system.
Delayed sends, idempotent enqueues, priority lanes, periodic digests, and batch fan-out for a notification system.
A notification service is a natural fit for a queue: every send is independent, some are scheduled for later, duplicates must be suppressed, and urgent alerts should jump the line. This example wires those requirements onto the Java SDK.
notifications/
Tasks.java # one task per channel + payload records
Service.java # the producer API the rest of the app calls
WorkerMain.java # the worker process
Each channel is its own task so it retries and scales independently.
import java.time.Duration;
import org.byteveda.flexiq.task.RetryPolicy;
import org.byteveda.flexiq.task.Task;
public final class Tasks {
public record Email(String to, String subject, String body) {}
public record Sms(String to, String text) {}
public static final Task<Email> SEND_EMAIL = Task.of("send_email", Email.class)
.maxRetries(5)
.retryPolicy(RetryPolicy.exponential(Duration.ofSeconds(2), Duration.ofMinutes(5)));
public static final Task<Sms> SEND_SMS = Task.of("send_sms", Sms.class)
.maxRetries(3);
public static final Task<String> SEND_DIGEST = Task.of("send_digest", String.class);
private Tasks() {}
}The producer-facing helpers map application events onto enqueue options.
import java.time.Duration;
import java.util.List;
import org.byteveda.flexiq.FlexiQ;
import org.byteveda.flexiq.model.JobFilter;
import org.byteveda.flexiq.model.JobStatus;
import org.byteveda.flexiq.task.EnqueueOptions;
public final class Service {
private final FlexiQ flexiq;
public Service(FlexiQ flexiq) {
this.flexiq = flexiq;
}
// 1. Delayed scheduling — remind 24h from now.
public String scheduleReminder(String email) {
return flexiq.enqueue(Tasks.SEND_EMAIL, new Tasks.Email(email, "Reminder", "..."),
EnqueueOptions.builder().delay(Duration.ofHours(24)).build());
}
// 2. Idempotency — a duplicate enqueue with the same key is a no-op while
// the first is still pending or running.
public String sendWelcomeOnce(String userId, String email) {
return flexiq.enqueue(Tasks.SEND_EMAIL, new Tasks.Email(email, "Welcome", "..."),
EnqueueOptions.builder().uniqueKey("welcome:" + userId).build());
}
// 3. Priority — security alerts preempt routine mail.
public String sendSecurityAlert(String phone, String text) {
return flexiq.enqueue(Tasks.SEND_SMS, new Tasks.Sms(phone, text),
EnqueueOptions.builder().priority(100).build());
}
// 4. Batch fan-out — stage a whole digest run in one storage call.
public List<String> sendDigests(List<Tasks.Email> due) {
return flexiq.enqueueMany(Tasks.SEND_EMAIL, due);
}
// 5. Cancellation — pull a not-yet-started send.
public boolean cancelScheduled(String jobId) {
return flexiq.cancel(jobId); // false if it already started
}
// 6. Inspection — what's still pending?
public List<org.byteveda.flexiq.model.Job> pendingSends() {
return flexiq.listJobs(JobFilter.builder()
.status(JobStatus.PENDING)
.task("send_email")
.limit(50)
.build());
}
}The worker runs the channels and registers the periodic digest at 08:00
America/New_York time.
import java.util.List;
import org.byteveda.flexiq.FlexiQ;
import org.byteveda.flexiq.scheduling.PeriodicTask;
import org.byteveda.flexiq.worker.Worker;
public final class WorkerMain {
public static void main(String[] args) throws InterruptedException {
try (FlexiQ flexiq = FlexiQ.builder().sqlite("notifications.db").open()) {
flexiq.registerPeriodic(PeriodicTask.builder("daily-digest", "send_digest", "0 8 * * *")
.timezone("America/New_York")
.build());
try (Worker worker = flexiq.worker()
.handle(Tasks.SEND_EMAIL, WorkerMain::deliverEmail)
.handle(Tasks.SEND_SMS, WorkerMain::deliverSms)
.handle(Tasks.SEND_DIGEST, ignored -> buildAndFanOutDigests(flexiq))
.queues("default")
.start()) {
worker.awaitShutdown();
}
}
}
private static Void deliverEmail(Tasks.Email email) {
// Call your email provider here.
System.out.println("email -> " + email.to());
return null;
}
private static Void deliverSms(Tasks.Sms sms) {
// Call your SMS provider here.
System.out.println("sms -> " + sms.to());
return null;
}
private static Void buildAndFanOutDigests(FlexiQ flexiq) {
// Build the due digest emails from your data store, then fan out.
List<Tasks.Email> due = List.of();
flexiq.enqueueMany(Tasks.SEND_EMAIL, due);
return null;
}
}java -cp app.jar WorkerMainjshell> new Service(flexiq).sendSecurityAlert("+1...", "Login from a new device")| Pattern | Where |
|---|---|
| Delayed / scheduled send | EnqueueOptions.builder().delay(...) |
| Idempotent send | EnqueueOptions.builder().uniqueKey(...) |
| Priority lane | EnqueueOptions.builder().priority(...) |
| Cancel a pending job | flexiq.cancel |
| Recurring digest | registerPeriodic (cron + timezone) |
| Batch fan-out | flexiq.enqueueMany |
| Pending-work inspection | flexiq.listJobs(JobFilter) |