nzw921rx commented on code in PR #11537:
URL: https://github.com/apache/seatunnel/pull/11537#discussion_r3635184907


##########
seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java:
##########
@@ -725,7 +725,33 @@ void aggregatedFlush(Map<SinkIdentifier, SinkContextProxy> 
proxyContexts) throws
                 for (SinkIdentifier id : sinkWritersWithIndex.get(i).keySet()) 
{
                     SinkContextProxy proxy = proxyContexts.get(id);
                     if (proxy != null && proxy.getFlushAction() != null) {
-                        proxy.getFlushAction().run();
+                        try {
+                            executeWithTableRetry(
+                                    id.getTableIdentifier(),
+                                    MultiTableFailurePhase.RUNTIME_WRITE,

Review Comment:
   Thank you for your contribution.🚀
   
   1. A new failure type should be introduced here; otherwise, it will be 
impossible to distinguish a `write` failure from a `flush` failure.
   
   2. I also have one concern. A flush operation handles a batch of records. On 
the JDBC side, the following logic is used to ensure that records are not 
inserted repeatedly during retries:
   
   ```java
   public void addToBatch(SeaTunnelRow record) throws SQLException {
       boolean exist = existRow(record);
       if (exist) {
           if (preExistFlag != null && !preExistFlag) {
               insertStatement.executeBatch();
               insertStatement.clearBatch();
           }
           rowConverter.toExternal(
                   valueTableSchema, databaseTableSchema, record, 
updateStatement);
           updateStatement.addBatch();
       } else {
           if (preExistFlag != null && preExistFlag) {
               updateStatement.executeBatch();
               updateStatement.clearBatch();
           }
           rowConverter.toExternal(
                   valueTableSchema, databaseTableSchema, record, 
insertStatement);
           insertStatement.addBatch();
       }
   
       preExistFlag = exist;
       submitted = false;
   }
   
   @Override
   public void executeBatch() throws SQLException {
       if (preExistFlag != null) {
           if (preExistFlag) {
               updateStatement.executeBatch();
               updateStatement.clearBatch();
           } else {
               insertStatement.executeBatch();
               insertStatement.clearBatch();
           }
       }
       submitted = true;
   }
   ```
   Do other connectors provide a similar mechanism for handling flush retries? 
If not, retrying a failed flush may cause unpredictable data issues, such as 
duplicate writes or inconsistent data.



-- 
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