diff --git a/conf/log4j2.properties b/conf/log4j2.properties index 94c52b38..707cdc5d 100644 --- a/conf/log4j2.properties +++ b/conf/log4j2.properties @@ -21,6 +21,8 @@ name = AccumuloTestingDefaultLoggingProperties status = info dest = err monitorInterval = 30 +# let ctrl-c logging keep working while ingest finishes its current flush +shutdownHook = disable appenders = console appender.console.type = Console diff --git a/src/main/java/org/apache/accumulo/testing/continuous/ContinuousIngest.java b/src/main/java/org/apache/accumulo/testing/continuous/ContinuousIngest.java index 62cad5da..4b671059 100644 --- a/src/main/java/org/apache/accumulo/testing/continuous/ContinuousIngest.java +++ b/src/main/java/org/apache/accumulo/testing/continuous/ContinuousIngest.java @@ -30,6 +30,7 @@ import java.util.SortedSet; import java.util.TreeSet; import java.util.UUID; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.function.LongSupplier; import java.util.zip.CRC32; @@ -60,6 +61,9 @@ public class ContinuousIngest { private static final byte[] EMPTY_BYTES = new byte[0]; + // how long ctrl-c waits for ingest to reach a flush point before giving up on a clean stop + private static final long STOP_WAIT_SEC = 300; + private static List visibilities; private static long lastPauseNs; private static long pauseWaitSec; @@ -74,6 +78,10 @@ public class ContinuousIngest { private static RandomDataGenerator rnd; + // set by the shutdown hook, causes ingest to stop at the next flush point + private static volatile boolean stopping = false; + private static final CountDownLatch stopped = new CountDownLatch(1); + public interface RandomGeneratorFactory extends Supplier { static RandomGeneratorFactory create(ContinuousEnv env, AccumuloClient client, Supplier> splitSupplier, Random random) { @@ -270,11 +278,28 @@ public static void main(String[] args) throws Exception { final boolean checksum = Boolean.parseBoolean(testProps.getProperty(TestProps.CI_INGEST_CHECKSUM)); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + if (stopped.getCount() == 0) + return; + stopping = true; + log.info("Stopping ingest at next flush point, waiting up to {}s (kill -9 {} to stop now)", + STOP_WAIT_SEC, ProcessHandle.current().pid()); + try { + if (!stopped.await(STOP_WAIT_SEC, TimeUnit.SECONDS)) { + log.warn("Timed out waiting for ingest to stop, exiting with data unflushed"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + })); + var splitSupplier = createSplitSupplier(client, tableName); var randomFactory = RandomGeneratorFactory.create(env, client, splitSupplier, random); var batchWriterFactory = BatchWriterFactory.create(client, env, splitSupplier); doIngest(client, randomFactory, batchWriterFactory, tableName, testProps, maxColF, maxColQ, numEntries, checksum, random); + } finally { + stopped.countDown(); } } @@ -365,7 +390,7 @@ protected static void doIngest(AccumuloClient client, RandomGeneratorFactory ran } lastFlushTime = flush(bw, entriesWritten, entriesDeleted, lastFlushTime); - if (entriesWritten >= numEntries) + if (entriesWritten >= numEntries || stopping) break out; pauseCheck(random); } @@ -401,7 +426,7 @@ protected static void doIngest(AccumuloClient client, RandomGeneratorFactory ran lastFlushTime = flush(bw, entriesWritten, entriesDeleted, lastFlushTime); } - if (entriesWritten >= numEntries) + if (entriesWritten >= numEntries || stopping) break out; pauseCheck(random); }