-
Notifications
You must be signed in to change notification settings - Fork 19
Expand file tree
/
Copy pathInMemoryQueueService.java
More file actions
executable file
·84 lines (72 loc) · 2.53 KB
/
InMemoryQueueService.java
File metadata and controls
executable file
·84 lines (72 loc) · 2.53 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
package com.example;
import java.io.IOException;
import java.io.InputStream;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.TimeUnit;
public class InMemoryQueueService implements PriorityQueueService {
private final Map<String, PriorityBlockingQueue<PriorityMessage>> queues;
private long visibilityTimeout;
InMemoryQueueService() {
this.queues = new ConcurrentHashMap<>();
String propFileName = "config.properties";
Properties confInfo = new Properties();
try (InputStream inStream = getClass().getClassLoader().getResourceAsStream(propFileName)) {
confInfo.load(inStream);
} catch (IOException e) {
e.printStackTrace();
}
this.visibilityTimeout = Integer.parseInt(confInfo.getProperty("visibilityTimeout", "30"));
}
@Override
public void push(String queueUrl, String msgBody, int rank) {
// priority --> numerically lowest rank, lowest time, message
PriorityBlockingQueue<PriorityMessage> queue = queues.get(queueUrl);
if (queue == null) {
queue = new PriorityBlockingQueue<PriorityMessage>();
queues.put(queueUrl, queue);
}
Long now = now();
PriorityMessage priorityMessage = new PriorityMessage(rank,now, new Message(msgBody));
queue.add(priorityMessage);
}
@Override
public Message pull(String queueUrl) {
PriorityBlockingQueue<PriorityMessage> queue = queues.get(queueUrl);
if (queue == null) {
return null;
}
long nowTime = now();
while (!queue.isEmpty()) {
PriorityMessage priorityMessage = queue.poll();
if (priorityMessage == null) {
return null;
} else {
Message msg = priorityMessage.getMessage();
msg.setReceiptId(UUID.randomUUID().toString());
msg.incrementAttempts();
msg.setVisibleFrom(nowTime + TimeUnit.SECONDS.toMillis(visibilityTimeout));
return new Message(msg.getBody(), msg.getReceiptId());
}
}
return null;
}
@Override
public void delete(String queueUrl, String receiptId) {
PriorityBlockingQueue<PriorityMessage> queue = queues.get(queueUrl);
if (queue != null) {
long nowTime = now();
for (PriorityMessage priorityMessage : queue) {
Message message = priorityMessage.getMessage();
if (!message.isVisibleAt(nowTime) && message.getReceiptId().equals(receiptId)) {
queue.remove(priorityMessage);
break;
}
}
}
}
long now() {
return System.currentTimeMillis();
}
}