Проверка общей коллекции в многопоточном приложении

У меня есть коллекция, в которой лежат объекты класса "Магазин". У каждого объекта "Магазин" есть числовое поле "Количество товаров". Потоки перемещают товары из одного магазина в другой. Я хочу создать один поток, который будет следить за тем, что общее число товаров в магазинах не изменилось. С синхронизацией процесса перемещения товаров проблем нет. Но для проверки я не могу синхронизироваться по общей коллекции, потому что это остановит работу всех остальных потоков. Но если этого не делать, пока проверяющий поток будет ходить по списку магазинов и суммировать число товаров, какой-то из потоков может переместить товар из непроверенного магазина в уже проверенный, и таким образом получится недостача. Как мне этого избежать? Не прошу писать код, просто подайте, пожалуйста, какую-нибудь идею. Я пробовал использовать CopyOnWriteArrayList, но он не дает нужного результата, потому что я получаю копию листа, а сами объекты, которые лежат внутри него, в это время все равно могут изменяться. Даже если я буду заменять объекты в листе на новые, при перемещении товаров из магазина в магазин будет совершаться две операции (обновление числа товаров в одном магазине и в другом), и в каждом случае будет создаваться копия, и я могу взять промежуточную версию.


Ответы (1 шт):

Автор решения: Roman Konoval

Идея с проверкой перемещений правильная и сработает. Нужно только реализовать ее так, чтоб не было блокировок потоков, которые совершают транзакции, т.е. которые перемещают данные.

Общая идея решения это создать журнал транзакций, в который каждая транзакция пишет изменение, которое она сделала. Поток-наблюдатель читает журнал операций и проверяет, что после каждой операции условие выполняется. Есть много нюансов как реализовать детали, я не буду во все вдаваться, но буду по ходу упоминать.

class Shop {
  Integer id;
  ...
}

class Transaction {
  // хранит для каждого идентификатора магазина изменение в количестве товара
  Map<Integer, Integer> changes;
}

ConcurrentLinkedQueue<Transaction> transactionLog = new ConcurrentLinkedQueue<>();

Теперь каждая транзакция должна добавлять в transactionLog изменение, которое она делает.

Чтобы избежать ошибок и злонамеренных действий, можно изменить API вашего хранилища магазинов так, чтоб к самим магазинам клиенты не имели доступа. А работали чисто с объектами транзакций, т.е. API хранилища будет не List<Shop> какой-нибудь, где можно достать любой магазин и сделать что угодно, а что-то вроде:

class Shops {
  private List<Shop> shops;

  // клиенты, которые хотят переместить что-то из магазина в магазин могут
  // делать это используя только этот метод 
  // список shops недоступен извне, а объекты Shop клиенты либо не могут достать вообще
  // либо они immutable, т.е. клиент не может изменять количество в обход process
  public void process(Transaction transaction) {
    // тут должна быть потокобезопасная обработка изменения
    // т.е. модификация значений в колекции shops используя данные из transaction
    ...
    
    transactionLog.add(transaction);
  }
}

Shops shops = ...

Тогда поток, который делает перемещение, будет выглядеть как-то так (тут я привожу вариант з неизменяемым(immutable) классом Shop):

// эта операция одновременно захватывает блокировку на запись к магазину
// чтоб никакой другой поток не мог параллельно его изменить
Shop shop1 = shops.getAndLock(SHOP_ID1);
Shop shop2 = shops.getAndLock(SHOP_ID2);
Map<Integer, Integer> moveTx = new HashMap<>();
moveTx.put(shop1.getId(), +20);
moveTx.put(shop2.getId(), -20);
shops.process(new Transaction(moveTx));
shops.releaseLock(shop2);
shops.releaseLock(shop1);

Поток-демон должен просто читать transactionLog и проверять, что после каждой транзакции у нас условие целостности выполняется.

Что касается блокировок/разблокировок отдельных магазинов в потоках. Пример, который я привел выше, демонстрирует, что должно происходить, само API конечно можно и нужно улучшить, чтобы не забыть разблокировать, например, и т.д.

Я имею ввиду что можно улучшить API класса Shops так:

  shops.update(asList(SHOP_ID1, SHOP_ID2), (Collection<Shop> shops, Transaction transaction) -> {
     transaction.add(SHOP_ID1, +20);
     transaction.add(SHOP_ID2, -20);
  });

  ...

  class Shops {
     
     interface TransactionCallback {
        void execute(Collection<Shop> shops, Transaction transaction);
     }
     
     void update(List<Integer> shopIds, TransactionCallback callback) {
       
       lockShops(shopIds);
       try {
          Transaction tx = new Transaction();
          List<Shop> shopsInTransaction = getShops(shopIds);
          callback.execute(shopsInTransaction, tx);
          process(tx); // теперь process будет private методом
       } finally {
         unlockShops(shopIds);
       } 
     }
  }
→ Ссылка