(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo, boolean missfire)
| 1097 | } |
| 1098 | |
| 1099 | private void fireCron(long timerSerialId, int serverId, @NotNull String timerId, long concurrentSerialNo, |
| 1100 | boolean missfire) { |
| 1101 | if (Task.call(zeze.newProcedure(() -> { |
| 1102 | var index = _tIndexs.get(timerId); |
| 1103 | if (index == null |
| 1104 | || index.getServerId() != zeze.getConfig().getServerId() // 不是拥有者,取消本地调度,应该是不大可能发生的。 |
| 1105 | || index.getSerialId() != timerSerialId // 新注册的,旧的future需要取消。 |
| 1106 | ) { |
| 1107 | Transaction.whileCommit(() -> cancelFuture(timerId)); |
| 1108 | return 0; |
| 1109 | } |
| 1110 | |
| 1111 | var nodeId = index.getNodeId(); |
| 1112 | var node = _tNodes.get(nodeId); |
| 1113 | if (node == null) { |
| 1114 | cancel(serverId, timerId, nodeId, null, null); // maybe concurrent cancel |
| 1115 | return 0; |
| 1116 | } |
| 1117 | |
| 1118 | var timer = node.getTimers().get(timerId); |
| 1119 | if (timer == null) |
| 1120 | throw new IllegalStateException("maybe operate before timer created"); |
| 1121 | var handle = findTimerHandle(timer.getHandleName()); |
| 1122 | var cronTimer = timer.getTimerObj_Zeze_Builtin_Timer_BCronTimer(); |
| 1123 | if (concurrentSerialNo == timer.getConcurrentFireSerialNo()) { |
| 1124 | var hasNext = Timer.nextCronTimer(cronTimer, missfire); |
| 1125 | var context = new TimerContext(this, timer, cronTimer.getHappenTimes(), |
| 1126 | cronTimer.getExpectedTime(), cronTimer.getNextExpectedTime()); |
| 1127 | |
| 1128 | // 当调度发生了错误或者由于异步时序没有原子保证,导致同时(或某个瞬间)在多个Server进程调度时, |
| 1129 | // 这个系列号保证触发用户回调只会发生一次。这个并发问题不取消定时器,继续尝试调度(去争抢执行权)。 |
| 1130 | // 定时器的调度生命期由其他地方保证最终一致。如果保证发生了错误,将一致并发争抢执行权。 |
| 1131 | var serialSaved = index.getSerialId(); |
| 1132 | var ret = Task.call(zeze.newProcedure(() -> { |
| 1133 | handle.onTimer(context); |
| 1134 | return 0; |
| 1135 | }, "Timer.fireCronUser." + timer.getHandleName())); |
| 1136 | |
| 1137 | var indexNew = _tIndexs.get(timerId); |
| 1138 | if (indexNew == null || indexNew.getSerialId() != serialSaved) |
| 1139 | return 0; // canceled or new timer. |
| 1140 | |
| 1141 | if (ret == Procedure.Exception) { |
| 1142 | // 用户处理不允许异常,其他错误记录忽略,日志已经记录。 |
| 1143 | cancel(serverId, timerId, nodeId, node, handle); |
| 1144 | return 0; |
| 1145 | } |
| 1146 | timer.setConcurrentFireSerialNo(concurrentSerialNo + 1); |
| 1147 | |
| 1148 | if (!hasNext) { |
| 1149 | cancel(serverId, timerId, nodeId, node, handle); |
| 1150 | return 0; |
| 1151 | } |
| 1152 | } |
| 1153 | // else 发生了并发执行争抢,也需要再次进行本地调度。此时直接使用cronTimer中的值,不需要再次进行计算。 |
| 1154 | |
| 1155 | // continue period |
| 1156 | scheduleCronNext(timerSerialId, serverId, timerId, |
no test coverage detected