Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions conf/log4j2.properties
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
dlmarion marked this conversation as resolved.

appenders = console
appender.console.type = Console
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<ColumnVisibility> visibilities;
private static long lastPauseNs;
private static long pauseWaitSec;
Expand All @@ -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<LongSupplier> {
static RandomGeneratorFactory create(ContinuousEnv env, AccumuloClient client,
Supplier<SortedSet<Text>> splitSupplier, Random random) {
Expand Down Expand Up @@ -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();
}
}

Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand Down
Loading