This is an automated email from the ASF dual-hosted git repository. clebertsuconic pushed a commit to branch 2.25.x in repository https://gitbox.apache.org/repos/asf/activemq-artemis.git
commit 50048dbac6fcc581734e51411ae0dfbf178d02a0 Author: Clebert Suconic <[email protected]> AuthorDate: Mon Sep 12 16:27:08 2022 -0400 ARTEMIS-3973 Solving concurrent issue between Subscription.ack and moveNext under heavy CPU usage (cherry picked from commit c4145e92267f33aadd6eb1161d9e980c6a8444a2) --- .../paging/cursor/impl/PageSubscriptionImpl.java | 12 +-- .../core/paging/cursor/impl/ConcurrentAckTest.java | 110 +++++++++++++++++++++ 2 files changed, 116 insertions(+), 6 deletions(-) diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/cursor/impl/PageSubscriptionImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/cursor/impl/PageSubscriptionImpl.java index 29c25b0e8b..61ca2448df 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/cursor/impl/PageSubscriptionImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/cursor/impl/PageSubscriptionImpl.java @@ -904,7 +904,7 @@ public final class PageSubscriptionImpl implements PageSubscription { getPageInfo(pageNr); } - private PageCursorInfo getPageInfo(final PagePosition pos) { + PageCursorInfo getPageInfo(final PagePosition pos) { return getPageInfo(pos.getPageNr()); } @@ -1057,7 +1057,7 @@ public final class PageSubscriptionImpl implements PageSubscription { // expressions private final AtomicInteger confirmed = new AtomicInteger(0); - public boolean isAck(int messageNumber) { + public synchronized boolean isAck(int messageNumber) { return completePage != null || acks.get(messageNumber) != null; } @@ -1082,7 +1082,7 @@ public final class PageSubscriptionImpl implements PageSubscription { } } - private PageCursorInfo(final long pageId, final int numberOfMessages) { + PageCursorInfo(final long pageId, final int numberOfMessages) { if (numberOfMessages < 0) { throw new IllegalStateException("numberOfMessages = " + numberOfMessages + " instead of being >=0"); } @@ -1145,11 +1145,11 @@ public final class PageSubscriptionImpl implements PageSubscription { checkDone(); } - public boolean isRemoved(final int messageNr) { + public synchronized boolean isRemoved(final int messageNr) { return removedReferences.get(messageNr) != null; } - public void remove(final int messageNr) { + public synchronized void remove(final int messageNr) { if (logger.isTraceEnabled()) { logger.tracef("PageCursor Removing messageNr %s on page %s", messageNr, pageId); } @@ -1187,7 +1187,7 @@ public final class PageSubscriptionImpl implements PageSubscription { } } - private boolean internalAddACK(final PagePosition position) { + synchronized boolean internalAddACK(final PagePosition position) { removedReferences.put(position.getMessageNr(), DUMMY); return acks.put(position.getMessageNr(), position) == null; } diff --git a/artemis-server/src/test/java/org/apache/activemq/artemis/core/paging/cursor/impl/ConcurrentAckTest.java b/artemis-server/src/test/java/org/apache/activemq/artemis/core/paging/cursor/impl/ConcurrentAckTest.java new file mode 100644 index 0000000000..1054369b6e --- /dev/null +++ b/artemis-server/src/test/java/org/apache/activemq/artemis/core/paging/cursor/impl/ConcurrentAckTest.java @@ -0,0 +1,110 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * <p> + * http://www.apache.org/licenses/LICENSE-2.0 + * <p> + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.activemq.artemis.core.paging.cursor.impl; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.activemq.artemis.api.core.SimpleString; +import org.apache.activemq.artemis.core.io.SequentialFileFactory; +import org.apache.activemq.artemis.core.paging.PagingManager; +import org.apache.activemq.artemis.core.paging.PagingStoreFactory; +import org.apache.activemq.artemis.core.paging.impl.PagingStoreImpl; +import org.apache.activemq.artemis.core.persistence.impl.nullpm.NullStorageManager; +import org.apache.activemq.artemis.core.settings.impl.AddressSettings; +import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; +import org.apache.activemq.artemis.utils.actors.ArtemisExecutor; +import org.junit.Assert; +import org.junit.Test; +import org.mockito.Mockito; + +public class ConcurrentAckTest extends ActiveMQTestBase { + + @Test + public void testConcurrentAddAckPaging() throws Throwable { + + ScheduledExecutorService scheduledExecutorService = Executors.newScheduledThreadPool(1); + runAfter(scheduledExecutorService::shutdownNow); + ExecutorService service = Executors.newFixedThreadPool(10); + runAfter(service::shutdownNow); + + for (int repeat = 0; repeat < 100; repeat++) { + // I needed brute force to make this test to fail, + // hence I am executing this method 100 times. + testConcurrentAddAckPaging(scheduledExecutorService, service); + } + } + + private void testConcurrentAddAckPaging(ScheduledExecutorService scheduledExecutorService, ExecutorService service) throws Throwable { + AtomicInteger errors = new AtomicInteger(0); + PagingStoreImpl store = new PagingStoreImpl(SimpleString.toSimpleString("TEST"), scheduledExecutorService, 100L, Mockito.mock(PagingManager.class), new NullStorageManager(), Mockito.mock(SequentialFileFactory.class), Mockito.mock(PagingStoreFactory.class), SimpleString.toSimpleString("TEST"), new AddressSettings(), ArtemisExecutor.delegate(service), ArtemisExecutor.delegate(service), false); + + PageCursorProviderImpl pageCursorProvider = new PageCursorProviderImpl(store, new NullStorageManager()); + PageSubscriptionImpl subscription = (PageSubscriptionImpl) pageCursorProvider.createSubscription(1, null, true); + PageSubscriptionImpl.PageCursorInfo cursorInfo = subscription.getPageInfo(new PagePositionImpl(1, 1)); + CountDownLatch done = new CountDownLatch(5); + + CyclicBarrier barrier = new CyclicBarrier(5); + + for (int r = 0; r < 4; r++) { + service.execute(() -> { + try { + barrier.await(1, TimeUnit.SECONDS); + } catch (Exception ignored) { + } + for (int i = 0; i < 5000; i++) { + try { + cursorInfo.internalAddACK(new PagePositionImpl(i, i)); + } catch (Throwable e) { + e.printStackTrace(); + errors.incrementAndGet(); + } + } + done.countDown(); + }); + } + + service.execute(() -> { + try { + try { + barrier.await(1, TimeUnit.SECONDS); + } catch (Exception ignored) { + } + for (int i = 0; i < 5000; i++) { + cursorInfo.isAck(i); + cursorInfo.isRemoved(i); + } + } catch (Exception e) { + e.printStackTrace(); + errors.incrementAndGet(); + } + + done.countDown(); + }); + + Assert.assertTrue(done.await(10, TimeUnit.SECONDS)); + + Assert.assertEquals(0, errors.get()); + } + +} \ No newline at end of file
