Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ public void start() {

@Override
public void close() {
super.close();
log.info("AlertHAServer shutdown...");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,11 @@ public interface ITaskGroupCoordinator extends AutoCloseable {
*/
void start();

/**
* Request the worker to stop without joining the worker. Call {@link #close()} to wait before restarting.
*/
void requestStop();

/**
* If the {@link TaskInstance#getTaskGroupId()} > 0, and the TaskGroup flag is {@link Flag#YES} then the task instance need to use task group.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,11 @@ public interface IWorkflowSerialCoordinator extends AutoCloseable {

void start();

/**
* Request the worker to stop without joining the worker. Call {@link #close()} to wait before restarting.
*/
void requestStop();

@Override
void close();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,7 @@
@Component
public class MasterCoordinator extends AbstractHAServer {

private final ITaskGroupCoordinator taskGroupCoordinator;

private final IFailoverCoordinator failoverCoordinator;

private final IWorkflowSerialCoordinator workflowSerialCoordinator;
private final MasterCoordinatorListener masterCoordinatorListener;

public MasterCoordinator(final Registry registry,
final MasterConfig masterConfig,
Expand All @@ -56,11 +52,9 @@ public MasterCoordinator(final Registry registry,
registry,
RegistryNodeType.MASTER_COORDINATOR.getRegistryPath(),
masterConfig.getMasterAddress());
this.taskGroupCoordinator = taskGroupCoordinator;
this.failoverCoordinator = failoverCoordinator;
this.workflowSerialCoordinator = workflowSerialCoordinator;
addServerStatusChangeListener(
new MasterCoordinatorListener(taskGroupCoordinator, failoverCoordinator, workflowSerialCoordinator));
this.masterCoordinatorListener =
new MasterCoordinatorListener(taskGroupCoordinator, failoverCoordinator, workflowSerialCoordinator);
addServerStatusChangeListener(masterCoordinatorListener);
}

@Override
Expand All @@ -71,7 +65,8 @@ public void start() {

@Override
public void close() {
taskGroupCoordinator.close();
super.close();
masterCoordinatorListener.changeToStandBy();
log.info("MasterCoordinator shutdown...");
}

Expand Down Expand Up @@ -108,11 +103,14 @@ public void changeToActive() {

@Override
public void changeToStandBy() {
taskGroupCoordinator.close();
workflowSerialCoordinator.close();
// Stop both workers before waiting: either may be blocked in a database call.
taskGroupCoordinator.requestStop();
workflowSerialCoordinator.requestStop();
if (failoverCoordinatorFuture != null) {
failoverCoordinatorFuture.cancel(true);
}
taskGroupCoordinator.close();
workflowSerialCoordinator.close();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ public class TaskGroupCoordinator implements ITaskGroupCoordinator, AutoCloseabl
@Autowired
private TransactionTemplate transactionTemplate;

private boolean flag = false;
private volatile boolean flag = false;

private Thread internalThread;

Expand Down Expand Up @@ -113,9 +113,19 @@ private void doStart() {
try {
final StopWatch taskGroupCoordinatorRoundCost = StopWatch.createStarted();

// A database call may outlive a stop request; check flag before starting the next phase.
amendTaskGroupUseSize();
if (!flag) {
return;
}
amendTaskGroupQueueStatus();
if (!flag) {
return;
}
dealWithForceStartTaskGroupQueue();
if (!flag) {
return;
}
dealWithWaitingTaskGroupQueue();

taskGroupCoordinatorRoundCost.stop();
Expand All @@ -124,7 +134,9 @@ private void doStart() {
log.error("TaskGroupCoordinator error", e);
} finally {
// sleep 5s
ThreadUtils.sleep(Constants.SLEEP_TIME_MILLIS * 5);
if (flag) {
ThreadUtils.sleep(Constants.SLEEP_TIME_MILLIS * 5);
}
}
}
}
Expand All @@ -141,7 +153,14 @@ private void amendTaskGroupUseSize() {
StopWatch taskGroupCoordinatorRoundTimeCost = StopWatch.createStarted();

for (TaskGroup taskGroup : taskGroups) {
if (!flag) {
return;
}
int actualUseSize = taskGroupQueueDao.countUsingTaskGroupQueueByGroupId(taskGroup.getId());
// The query may finish after a stop request; do not start the update in that case.
if (!flag) {
return;
}
if (taskGroup.getUseSize() == actualUseSize) {
continue;
}
Expand All @@ -161,10 +180,10 @@ private void amendTaskGroupQueueStatus() {
int minTaskGroupQueueId = -1;
int limit = DEFAULT_LIMIT;
StopWatch taskGroupCoordinatorRoundTimeCost = StopWatch.createStarted();
while (true) {
while (flag) {
List<TaskGroupQueue> taskGroupQueues =
taskGroupQueueDao.queryInQueueTaskGroupQueue(minTaskGroupQueueId, limit);
if (CollectionUtils.isEmpty(taskGroupQueues)) {
if (!flag || CollectionUtils.isEmpty(taskGroupQueues)) {
break;
}
amendTaskGroupQueueStatus(taskGroupQueues);
Expand All @@ -188,6 +207,9 @@ private void amendTaskGroupQueueStatus(List<TaskGroupQueue> taskGroupQueues) {
.collect(Collectors.toMap(TaskInstance::getId, Function.identity()));

for (TaskGroupQueue taskGroupQueue : taskGroupQueues) {
if (!flag) {
return;
}
int taskId = taskGroupQueue.getTaskId();
final TaskInstance taskInstance = taskInstanceMap.get(taskId);

Expand All @@ -214,10 +236,10 @@ private void dealWithForceStartTaskGroupQueue() {
int minTaskGroupQueueId = -1;
int limit = DEFAULT_LIMIT;
StopWatch taskGroupCoordinatorRoundTimeCost = StopWatch.createStarted();
while (true) {
while (flag) {
final List<TaskGroupQueue> taskGroupQueues =
taskGroupQueueDao.queryWaitNotifyForceStartTaskGroupQueue(minTaskGroupQueueId, limit);
if (CollectionUtils.isEmpty(taskGroupQueues)) {
if (!flag || CollectionUtils.isEmpty(taskGroupQueues)) {
break;
}
dealWithForceStartTaskGroupQueue(taskGroupQueues);
Expand All @@ -235,6 +257,9 @@ private void dealWithForceStartTaskGroupQueue(List<TaskGroupQueue> taskGroupQueu
// Notify the related waiting task instance
// Set the taskGroupQueue status to RELEASE and remove it from queue
for (final TaskGroupQueue taskGroupQueue : taskGroupQueues) {
if (!flag) {
return;
}
try {
LogUtils.setTaskInstanceIdMDC(taskGroupQueue.getTaskId());
if (!notifyForceStartTaskGroupQueue(taskGroupQueue)) {
Expand Down Expand Up @@ -290,6 +315,9 @@ private void dealWithWaitingTaskGroupQueue() {
return;
}
for (TaskGroup taskGroup : taskGroups) {
if (!flag) {
return;
}
int availableSize = taskGroup.getGroupSize() - taskGroup.getUseSize();
if (availableSize <= 0) {
log.info("TaskGroup {} is full, available size is {}", taskGroup, availableSize);
Expand All @@ -307,6 +335,9 @@ private void dealWithWaitingTaskGroupQueue() {
continue;
}
for (TaskGroupQueue taskGroupQueue : taskGroupQueues) {
if (!flag) {
return;
}
try {
LogUtils.setTaskInstanceIdMDC(taskGroupQueue.getTaskId());
if (!acquireTaskGroupSlotAndNotify(taskGroupQueue)) {
Expand Down Expand Up @@ -516,21 +547,36 @@ private void deleteTaskGroupQueueSlot(TaskGroupQueue taskGroupQueue) {
log.debug("TaskGroupQueue has already been released: {}", taskGroupQueue.getId());
}

@Override
public synchronized void requestStop() {
flag = false;
if (internalThread != null) {
internalThread.interrupt();
}
}

@Override
public synchronized void close() {
if (!flag) {
log.warn("TaskGroupCoordinator is already closed");
return;
if (Thread.currentThread() == internalThread) {
throw new IllegalStateException("TaskGroupCoordinator cannot close its own worker thread");
}
flag = false;
try {
if (internalThread != null) {
internalThread.interrupt();
// A prior stop request does not mean the worker has finished.
requestStop();
boolean interrupted = false;
if (internalThread != null) {
// Keep start() waiting until the old worker has finished, including any in-flight JDBC call.
while (internalThread.isAlive()) {
try {
internalThread.join();
} catch (InterruptedException ex) {
interrupted = true;
}
}
} catch (Exception ex) {
log.error("Close internalThread failed", ex);
internalThread = null;
}
if (interrupted) {
Thread.currentThread().interrupt();
}
internalThread = null;
log.info("TaskGroupCoordinator closed");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -95,15 +95,22 @@ private void doStart() {
try {
final StopWatch workflowSerialCoordinatorRoundCost = StopWatch.createStarted();
final List<SerialCommandsGroup> serialCommandsGroups = fetchSerialCommands();
serialCommandsGroups.forEach(this::handleSerialCommand);
// Fetching or handling a group may outlive a stop request; check before handling the next group.
serialCommandsGroups.forEach(serialCommandsGroup -> {
if (flag) {
handleSerialCommand(serialCommandsGroup);
}
});
log.debug("WorkflowSerialCoordinator handled SerialCommandsGroup size: {}, cost: {}/ms ",
serialCommandsGroups.size(),
workflowSerialCoordinatorRoundCost.getDuration().toMillis());
} catch (Throwable e) {
log.error("WorkflowSerialCoordinator error", e);
} finally {
// sleep 5s
ThreadUtils.sleep(TimeUnit.SECONDS.toMillis(DEFAULT_FETCH_INTERVAL_SECONDS));
if (flag) {
ThreadUtils.sleep(TimeUnit.SECONDS.toMillis(DEFAULT_FETCH_INTERVAL_SECONDS));
}
}
}
}
Expand Down Expand Up @@ -168,21 +175,36 @@ private SerialCommandsGroup createSerialCommandsGroup(SerialCommandDto serialCom
.build();
}

@Override
public synchronized void requestStop() {
flag = false;
if (internalThread != null) {
internalThread.interrupt();
}
}

@Override
public synchronized void close() {
if (!flag) {
log.warn("WorkflowSerialCoordinator is already closed");
return;
if (Thread.currentThread() == internalThread) {
throw new IllegalStateException("WorkflowSerialCoordinator cannot close its own worker thread");
}
flag = false;
try {
if (internalThread != null) {
internalThread.interrupt();
// A prior stop request does not mean the worker has finished.
requestStop();
boolean interrupted = false;
if (internalThread != null) {
// Keep start() waiting until the old worker has finished, including any in-flight JDBC call.
while (internalThread.isAlive()) {
try {
internalThread.join();
} catch (InterruptedException ex) {
interrupted = true;
}
}
} catch (Exception ex) {
log.error("Close internalThread failed", ex);
internalThread = null;
}
if (interrupted) {
Thread.currentThread().interrupt();
}
internalThread = null;
log.info("WorkflowSerialCoordinator closed");
}
}
Loading
Loading