قبل عامين ، واجهنا في أحد المشاريع الحاجة إلى تأجيل تنفيذ إجراء ما لفترة زمنية معينة. على سبيل المثال ، تعرف على حالة الدفع في غضون ثلاث ساعات أو أعد إرسال الإشعار بعد 45 دقيقة. ومع ذلك ، في ذلك الوقت لم نجد مكتبات مناسبة يمكنها "التأجيل" ولم تتطلب وقتًا إضافيًا للتهيئة والتشغيل. قمنا بتحليل الخيارات الممكنة وكتبنا مكتبة الانتظار الصغيرة الخاصة بنا في Java باستخدام Redis كمستودع. سأتحدث في هذا المقال عن إمكانيات المكتبة وبدائلها و "مكابس" التي وجدناها في هذه العملية.
وظائف
إذن ماذا تفعل قائمة الانتظار المتأخرة؟ يتم تسليم حدث مضاف إلى قائمة الانتظار المعلقة إلى المعالج بعد الفترة الزمنية المحددة. إذا فشلت المعالجة ، فسيتم تسليم الحدث مرة أخرى لاحقًا. علاوة على ذلك ، فإن الحد الأقصى لعدد المحاولات محدود. لا تضمن Redis السلامة ، وعليك أن تكون مستعدًا لخسارة الأحداث . ومع ذلك ، في إصدار الكتلة ، يُظهر Redis موثوقية عالية إلى حد ما ، ولم نواجه هذا مطلقًا خلال عام ونصف من التشغيل.
API
أضف حدثًا إلى قائمة الانتظار
eventService.enqueueWithDelayNonBlocking(new DummyEvent("1"), Duration.ofHours(1)).subscribe();
لاحظ أن الطريقة ترجع Mono، لذلك للتشغيل ، عليك القيام بأحد الإجراءات التالية:
subscribe(...)block()
يتم توفير تفسيرات أكثر تفصيلاً في الوثائق الخاصة بـ Project Reactor. يضاف السياق إلى الحدث مثل هذا:
eventService.enqueueWithDelayNonBlocking(new DummyEvent("1"), Duration.ofHours(1), Map.of("key", "value")).subscribe();
تسجيل معالج الحدث
eventService.addHandler(DummyEvent.class, e -> Mono.just(true), 1);
, :
eventService.addHandler(
DummyEvent.class,
e -> Mono
.subscriberContext()
.doOnNext(ctx -> {
Map<String, String> eventContext = ctx.get("eventContext");
log.info("context key {}", eventContext.get("key"));
})
.thenReturn(true),
1
);
eventService.removeHandler(DummyEvent.class);
"-":
import static com.github.fred84.queue.DelayedEventService.delayedEventService;
var eventService = delayedEventService().client(redisClient).build();
:
import static com.github.fred84.queue.DelayedEventService.delayedEventService;
var eventService = delayedEventService()
.client(redisClient)
.mapper(objectMapper)
.handlerScheduler(Schedulers.fromExecutorService(executor))
.schedulingInterval(Duration.ofSeconds(1))
.schedulingBatchSize(SCHEDULING_BATCH_SIZE)
.enableScheduling(false)
.pollingTimeout(POLLING_TIMEOUT)
.eventContextHandler(new DefaultEventContextHandler())
.dataSetPrefix("")
.retryAttempts(10)
.metrics(new NoopMetrics())
.refreshSubscriptionsInterval(Duration.ofMinutes(5))
.build();
( Redis) eventService.close() , @javax.annotation.PreDestroy.
- , . :
- , Redis;
- , ( "delayed.queue.ready.for.handling.count" )
, delayed queue. 2018
Amazon Web Services.
, . : " , Amazon-, ".
:
- , JMS . SQS , 15 .
" " . , Redis :
- sorted sets,
- "sorted_set" "list" ( )
, Netflix dyno-queues
. , , .
, " " sorted set list, ( ):
var events = redis.zrangebyscore("delayed_events", Range.create(-1, System.currentTimeMillis()), 100);
events.forEach(key -> {
var payload = extractPayload(key);
var listName = extractType(key);
redis.lpush(listName, payload);
redis.zrem("delayed_events", key);
});
redis.brpop(listName)
.
"list" (, ), list . Redis , 2 .
events.forEach(key -> {
...
redis.multi();
redis.zrem("delayed_events", key);
redis.lpush(listName, payload);
redis.exec();
});
list-a . , . "sorted_set" .
events.forEach(key -> {
...
redis.multi();
redis.zadd("delayed_events", nextAttempt(key))
redis.zrem("delayed_events", key);
redis.lpush(listName, payload);
redis.exec();
});
, , " " "delayed queue" . "sorted set"
metadata;payload, payload , metadata - . . , metadata payload Redis hset "sorted set" .
var envelope = metadata + SEPARATOR + payload;
redis.zadd(envelope, scheduledAt);
var envelope = metadata + SEPARATOR + payload;
var key = eventType + SEPARATOR + eventId;
redis.multi();
redis.zadd(key, scheduledAt);
redis.hset("metadata", key, envelope)
redis.exec();
, . , list . TTL :
redis.set(lockKey, "value", ex(lockTimeout.toMillis() * 1000).nx());
Spring, . " " :
Lettuce , . Project Reactor , " ".
, Subscriber
redis
.reactive()
.brpop(timeout, queue)
.map(e -> deserialize(e))
.subscribe(new InnerSubscriber<>(handler, ... params ..))
class InnerSubscriber<T extends Event> extends BaseSubscriber<EventEnvelope<T>> {
@Override
protected void hookOnNext(@NotNull EventEnvelope<T> envelope) {
Mono<Boolean> promise = handler.apply(envelope.getPayload());
promise.subscribe(r -> request(1));
}
}
, ( Netflix dyno queue, poll- ).
?
- Kotlin DSL. Kotlin
suspend funAPIProject Reactor