This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new fb2684a8249 Export TsFile: manual commit instead of auto commit 
(#15798)
fb2684a8249 is described below

commit fb2684a8249cea3eee5edbf95443b9d80f1b7e43
Author: VGalaxies <[email protected]>
AuthorDate: Mon Jun 23 17:10:29 2025 +0800

    Export TsFile: manual commit instead of auto commit (#15798)
---
 .../cli/src/main/java/org/apache/iotdb/tool/common/Constants.java      | 3 +--
 .../apache/iotdb/tool/tsfile/subscription/SubscriptionTableTsFile.java | 2 +-
 .../apache/iotdb/tool/tsfile/subscription/SubscriptionTreeTsFile.java  | 2 +-
 3 files changed, 3 insertions(+), 4 deletions(-)

diff --git 
a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/common/Constants.java 
b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/common/Constants.java
index c93308289e6..049291ac248 100644
--- a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/common/Constants.java
+++ b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/common/Constants.java
@@ -285,9 +285,8 @@ public class Constants {
   public static final String LOOSE_RANGE = "";
   public static final boolean STRICT = false;
   public static final String MODE = "snapshot";
-  public static final boolean AUTO_COMMIT = true;
+  public static final boolean AUTO_COMMIT = false;
   public static final String TABLE_MODEL = "table";
-  public static final long AUTO_COMMIT_INTERVAL = 5000;
   public static final long POLL_MESSAGE_TIMEOUT = 10000;
   public static final String TOPIC_NAME_PREFIX = "topic_";
   public static final String GROUP_NAME_PREFIX = "group_";
diff --git 
a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTableTsFile.java
 
b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTableTsFile.java
index f869490e7c5..3a8bee30b43 100644
--- 
a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTableTsFile.java
+++ 
b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTableTsFile.java
@@ -106,7 +106,6 @@ public class SubscriptionTableTsFile extends 
AbstractSubscriptionTsFile {
                   .consumerId(Constants.CONSUMER_NAME_PREFIX + i)
                   .consumerGroupId(groupId)
                   .autoCommit(Constants.AUTO_COMMIT)
-                  .autoCommitIntervalMs(Constants.AUTO_COMMIT_INTERVAL)
                   .fileSaveDir(commonParam.getTargetDir())
                   .build());
     }
@@ -163,6 +162,7 @@ public class SubscriptionTableTsFile extends 
AbstractSubscriptionTsFile {
                       throw new RuntimeException(e);
                     }
                     commonParam.getCountFile().incrementAndGet();
+                    consumer.commitSync(message);
                   }
                 } catch (Exception e) {
                   e.printStackTrace(System.out);
diff --git 
a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTreeTsFile.java
 
b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTreeTsFile.java
index 7f2d302e6e4..4d8b682aa22 100644
--- 
a/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTreeTsFile.java
+++ 
b/iotdb-client/cli/src/main/java/org/apache/iotdb/tool/tsfile/subscription/SubscriptionTreeTsFile.java
@@ -103,7 +103,6 @@ public class SubscriptionTreeTsFile extends 
AbstractSubscriptionTsFile {
                   .consumerId(Constants.CONSUMER_NAME_PREFIX + i)
                   .consumerGroupId(groupId)
                   .autoCommit(Constants.AUTO_COMMIT)
-                  .autoCommitIntervalMs(Constants.AUTO_COMMIT_INTERVAL)
                   .fileSaveDir(commonParam.getTargetDir())
                   .buildPullConsumer());
     }
@@ -159,6 +158,7 @@ public class SubscriptionTreeTsFile extends 
AbstractSubscriptionTsFile {
                       throw new RuntimeException(e);
                     }
                     commonParam.getCountFile().incrementAndGet();
+                    consumer.commitSync(message);
                   }
                 } catch (Exception e) {
                   e.printStackTrace(System.out);

Reply via email to