Rangsh commented on PR #12511:
URL: https://github.com/apache/seatunnel/pull/12511#issuecomment-5977952591

   @Vivek1106-04 @DanielLeens
   
   Thanks again for putting the complete fix together. On the blocker in 
DanielLeens's review (Issue 1: nothing fails if `|| savepointDraining` is 
removed): #12493 had a test for exactly that window, 
`testSavepointDrainGateReArmsWhenPendingCounterZero`, so I ported it onto this 
branch in case it saves you some time. I adapted it to the plain `boolean` 
field and the helpers here, and added `SCHEMA_CHANGE_AFTER_POINT_TYPE`.
   
   It sets `pendingCounter = 0` with the gate set. That is the gap between the 
in-flight checkpoint finishing and the drain thread waking up, where only the 
gate keeps a trigger from creating a checkpoint ahead of the savepoint. It 
asserts that every trigger type creates nothing and re-arms, and then that 
periodic triggering resumes once the gate is cleared.
   
   Checked locally at `91001c42c`:
   
   - with `|| savepointDraining` removed: `a CHECKPOINT_TYPE trigger created a 
checkpoint while the drain gate was set ==> expected: <0> but was: <1>`. The 
existing 
`testTriggerDuringSavepointDrainReturnsAndRearmsWithoutCreatingACheckpoint` 
still passes in the same run, which is the gap the review points out.
   - with the branch unchanged: `CheckpointCoordinatorTest` `Tests run: 21, 
Failures: 0, Errors: 0`.
   
   <details><summary>Test method (drop-in after 
<code>testTriggerDuringSavepointDrainReturnsAndRearmsWithoutCreatingACheckpoint</code>)</summary>
   
   ```java
   /**
    * Between the in-flight checkpoint finishing and the drain thread waking 
up, {@code
    * pendingCounter} is already 0 and only the drain gate keeps a trigger from 
creating a
    * checkpoint ahead of the savepoint.
    */
   @Test
   void testTriggerWhileSavepointDrainGateIsSetAndNothingIsPendingRearms() {
       ExecutorService executor = Executors.newCachedThreadPool();
       CheckpointCoordinator coordinator = 
Mockito.spy(buildMinimalCoordinator(executor));
       try {
           List<CheckpointType> rearmed = new CopyOnWriteArrayList<>();
           Mockito.doAnswer(
                           invocation -> {
                               rearmed.add(invocation.getArgument(0));
                               return null;
                           })
                   .when(coordinator)
                   .scheduleTriggerPendingCheckpoint(
                           Mockito.any(CheckpointType.class), 
Mockito.anyLong());
           Mockito.doReturn(new InvocationFuture[0])
                   .when(coordinator)
                   .triggerCheckpoint(Mockito.any(CheckpointBarrier.class));
           isAllTaskReady(coordinator).set(true);
           pendingCounter(coordinator).set(0);
           ReflectionUtils.setField(
                   coordinator, CheckpointCoordinator.class, 
"savepointDraining", true);
   
           for (CheckpointType type :
                   new CheckpointType[] {
                       CheckpointType.CHECKPOINT_TYPE,
                       CheckpointType.COMPLETED_POINT_TYPE,
                       CheckpointType.SCHEMA_CHANGE_BEFORE_POINT_TYPE,
                       CheckpointType.SCHEMA_CHANGE_AFTER_POINT_TYPE
                   }) {
               coordinator.tryTriggerPendingCheckpoint(type);
               Assertions.assertEquals(
                       0,
                       pendingCounter(coordinator).get(),
                       "a " + type + " trigger created a checkpoint while the 
drain gate was set");
               Assertions.assertTrue(
                       rearmed.contains(type),
                       "a " + type + " trigger was dropped instead of 
re-armed");
           }
   
           ReflectionUtils.setField(
                   coordinator, CheckpointCoordinator.class, 
"savepointDraining", false);
           assertPeriodicTriggeringResumes(coordinator);
       } finally {
           // no savepoint request backs the gate set above
           ReflectionUtils.setField(
                   coordinator, CheckpointCoordinator.class, 
"savepointDraining", false);
           shutDown(coordinator);
           executor.shutdownNow();
       }
   }
   ```
   
   </details>
   
   Please feel free to take it as is, rename it, or skip it entirely; it's your 
call. I'm not touching the other review items. If it would be easier, I'm happy 
to open a PR against your branch instead.
   
   On the Build: as DanielLeens noted in the review, the failed jobs are in 
modules this PR doesn't change. The only one of them that runs a savepoint is 
`PostgresCDCIT`, and in that run the savepoint and restore themselves succeeded 
(`SAVEPOINT_DONE`, then restore from checkpoint 4). The row missing after 
restore has the same signature as the open #12382, which was seen on a `dev` 
scheduled run before this PR existed (see also #11847). The branch is now 41 
commits behind `dev`, which includes #12444 for the PayPal test, so syncing 
`dev` in the same push and re-running should give a cleaner signal.
   
   Thanks!
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to