001/* 002 * (C) Copyright 2016 Nuxeo SA (http://nuxeo.com/) and others. 003 * 004 * Licensed under the Apache License, Version 2.0 (the "License"); 005 * you may not use this file except in compliance with the License. 006 * You may obtain a copy of the License at 007 * 008 * http://www.apache.org/licenses/LICENSE-2.0 009 * 010 * Unless required by applicable law or agreed to in writing, software 011 * distributed under the License is distributed on an "AS IS" BASIS, 012 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 013 * See the License for the specific language governing permissions and 014 * limitations under the License. 015 * 016 */ 017package org.nuxeo.ecm.platform.importer.queue.manager; 018 019import java.util.ArrayList; 020import java.util.List; 021import java.util.concurrent.ArrayBlockingQueue; 022import java.util.concurrent.BlockingQueue; 023import java.util.concurrent.TimeUnit; 024 025import org.nuxeo.ecm.platform.importer.log.ImporterLogger; 026import org.nuxeo.ecm.platform.importer.source.SourceNode; 027 028/** 029 * @since 8.3 030 */ 031public abstract class AbstractQueuesManager implements QueuesManager { 032 033 List<BlockingQueue<SourceNode>> queues = new ArrayList<BlockingQueue<SourceNode>>(); 034 035 protected int maxQueueSize = 1000; 036 037 protected ImporterLogger log = null; 038 039 public AbstractQueuesManager(ImporterLogger logger, int queuesNb, int maxQueueSize) { 040 this.maxQueueSize = maxQueueSize; 041 for (int i = 0; i < queuesNb; i++) { 042 queues.add(new ArrayBlockingQueue<SourceNode>(maxQueueSize)); 043 } 044 log = logger; 045 } 046 047 @Override 048 public BlockingQueue<SourceNode> getQueue(int idx) { 049 return queues.get(idx); 050 } 051 052 @Override 053 public boolean isQueueEmpty(int idQueue) { 054 return queues.get(idQueue).isEmpty(); 055 } 056 057 @Override 058 public int dispatch(SourceNode bh) throws InterruptedException { 059 int idx = getTargetQueue(bh, queues.size()); 060 061 boolean accepted = getQueue(idx).offer(bh, 1, TimeUnit.SECONDS); 062 063 if (!accepted) { 064 log.warn("Timeout while waiting for an available queue"); 065 idx = getTargetQueue(bh, queues.size()); 066 getQueue(idx).offer(bh, 5, TimeUnit.SECONDS); 067 } 068 return idx; 069 } 070 071 protected abstract int getTargetQueue(SourceNode bh, int nbQueues); 072 073 @Override 074 public int getNBConsumers() { 075 return queues.size(); 076 } 077}