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}