我试图写它有两个方法批量邮件服务:生产者/用料和冲洗功能的消费者
add(Mail mail)
:可发送Email,由生产者
flushMailService()
称为:刷新服务。消费者应该列出一个清单,并打电话给另一个(昂贵的)方法。通常只有在达到批量大小后才能调用昂贵的方法。
这是可以做到与poll()
具有超时。但是生产者应该能够清除邮件服务,如果它不想等待超时,而是让生产者发送队列中的邮件。
poll(20, TimeUnit.SECONDS)
可以被中断。如果发生中断,则应发送队列中的所有邮件,而不管批量大小是否已达到队列为空(使用poll()
,如果队列为空则立即返回null
;一旦该队列为空,则发送的邮件该中断已经被送到制片人。然后,生产者应当再次致电然后阻断poll
版本,直到其他任何生产中断等。
这似乎与给定的实施工作。
我试着使用ExecutorServices与Futures,但它似乎只能被中断一次,因为它们是在第一次中断后视为取消。因此,我使用了可以多次中断的线程。
目前我有以下实现似乎工作(但正在使用“原始”线程)。
这是一个合理的方法吗?或者也许可以使用另一种方法?
public class BatchMailService {
private LinkedBlockingQueue<Mail> queue = new LinkedBlockingQueue<>();
private CopyOnWriteArrayList<Thread> threads = new CopyOnWriteArrayList<>();
private static Logger LOGGER = LoggerFactory.getLogger(BatchMailService.class);
public void checkMails() {
int batchSize = 100;
int timeout = 20;
int consumerCount = 5;
Runnable runnable =() -> {
boolean wasInterrupted = false;
while (true) {
List<Mail> buffer = new ArrayList<>();
while (buffer.size() < batchSize) {
try {
Mail mail;
wasInterrupted |= Thread.interrupted();
if (wasInterrupted) {
mail = queue.poll(); // non-blocking call
} else {
mail = queue.poll(timeout, TimeUnit.SECONDS); // blocking call
}
if (mail != null) { // mail found immediately, or within timeout
buffer.add(mail);
} else { // no mail in queue, or timeout reached
LOGGER.debug("{} all mails currently in queue have been processed", Thread.currentThread());
wasInterrupted = false;
break;
}
} catch (InterruptedException e) {
LOGGER.info("{} interrupted", Thread.currentThread());
wasInterrupted = true;
break;
}
}
if (!buffer.isEmpty()) {
LOGGER.info("{} sending {} mails", Thread.currentThread(), buffer.size());
mailService.sendMails(buffer);
}
}
};
LOGGER.info("starting 5 threads ");
for (int i = 0; i < 5; i++) {
Thread thread = new Thread(runnable);
threads.add(thread);
thread.start();
}
}
public void addMail(Mail mail) {
queue.add(mail);
}
public void flushMailService() {
LOGGER.info("flushing BatchMailService");
for (Thread t : threads) {
t.interrupt();
}
}
}
无中断的另一种方法,但毒丸计划(Mail POISON_PILL = new Mail()
)的变体可能是以下。当有一个消费者线程时可能效果最好。至少,对于一个毒丸,只有一个消费者会继续。
Runnable runnable =() -> {
boolean flush = false;
boolean shutdown = false;
while (!shutdown) {
List<Mail> buffer = new ArrayList<>();
while (buffer.size() < batchSize && !shutdown) {
try {
Mail mail;
if (flush){
mail = queue.poll();
if (mail == null) {
LOGGER.info(Thread.currentThread() + " all mails currently in queue have been processed");
flush = false;
break;
}
}else {
mail = queue.poll(5, TimeUnit.SECONDS); // blocking call
}
if (mail == POISON_PILL){ // flush
LOGGER.info(Thread.currentThread() + " got flush");
flush = true;
}
else if (mail != null){
buffer.add(mail);
}
} catch (InterruptedException e) {
LOGGER.info(Thread.currentThread() + " interrupted");
shutdown = true;
}
}
if (!buffer.isEmpty()) {
LOGGER.info(Thread.currentThread()+"{} sending " + buffer.size()+" mails");
mailService.sendEmails(buffer);
}
}
};
public void flushMailService() {
LOGGER.info("flushing BatchMailService");
queue.add(POISON_PILL);
}
那么,如果生产者检查队列和刷新的尺寸达到当批量大小,然后如果你有多个消费者,可能他们不会用更少的邮件的邮件服务多次打电话? – user140547
它会有但锁定,以防止这种情况发生。可能有一个错误,我错过了,但正如我看到它正常工作。 addMail()在第二次检查queue.size()之前获取'fluskLock'。 当'checkMails()'从'await()'唤醒时,它会在继续刷新之前获取'flushLock'。 [链接](https://docs.oracle.com/javase/7/docs/api/java/util/concurrent/locks/Condition。html#awaitNanos(long)) –
我用5个并行生产者和1500邮件运行代码。在没有重复邮件的情况下,没有任何批次的BATCH_SIZE项目少于其中的项目,尽管它们通常有更多。 –