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

HTHou 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 d6b3afbb2b6 Add audit hooks for inter-node user data transfers (#18569)
d6b3afbb2b6 is described below

commit d6b3afbb2b67d251e27d1239cae5b5d833b61b7b
Author: Haonan <[email protected]>
AuthorDate: Thu Sep 3 14:35:54 2026 +0800

    Add audit hooks for inter-node user data transfers (#18569)
    
    * [Feature] Add audit hooks for inter-node user data transfers
    
    * [Feature] Address review for user data transfer audit hooks
    
    * [Feature] Address second-pass audit review
    
    * [Feature] Gate audit classification and filter snapshots
    
    * [Feature] Address audit transfer review feedback
    
    * [Feature] Reduce disabled transfer audit overhead
    
    * [Feature] Exclude audit database from IoTConsensusV2 audit
    
    * [Feature] Avoid duplicate IoTConsensusV2 audit classification
---
 .../iotdb/consensus/i18n/IoTConsensusMessages.java |   3 +
 .../iotdb/consensus/i18n/IoTConsensusMessages.java |   3 +
 .../common/request/IndexedConsensusRequest.java    |  10 ++
 .../iotdb/consensus/config/ConsensusConfig.java    |  41 ++++-
 .../config/UserDataTransferAuditClassifier.java    |  47 ++++++
 .../apache/iotdb/consensus/iot/IoTConsensus.java   |  14 +-
 .../consensus/iot/IoTConsensusServerImpl.java      | 124 ++++++++++++++-
 .../consensus/iot/client/DispatchLogHandler.java   | 131 ++++++++++++++--
 .../iotdb/consensus/iot/logdispatcher/Batch.java   |  10 ++
 .../consensus/iot/logdispatcher/LogDispatcher.java |   8 +-
 .../consensus/iot/IoTConsensusServerImplTest.java  |  36 +++++
 .../iot/client/DispatchLogHandlerTest.java         | 173 +++++++++++++++++++++
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |   2 +
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |   2 +
 .../db/audit/DataNodeUserDataTransferAuditor.java  | 135 ++++++++++++++++
 .../db/consensus/DataRegionConsensusImpl.java      |  24 +++
 .../iotconsensusv2/IoTConsensusV2AsyncSink.java    |  46 +++++-
 .../iotconsensusv2/IoTConsensusV2SyncSink.java     |  66 ++++++++
 .../IoTConsensusV2TabletBatchEventHandler.java     |   8 +
 .../IoTConsensusV2TabletInsertionEventHandler.java |  11 ++
 .../IoTConsensusV2TsFileInsertionEventHandler.java |  16 ++
 .../execution/exchange/source/SourceHandle.java    |  63 ++++++--
 .../plan/scheduler/AsyncPlanNodeSender.java        |  46 +++++-
 .../plan/scheduler/AsyncSendPlanNodeHandler.java   |  53 ++++++-
 .../scheduler/FragmentInstanceDispatcherImpl.java  |  45 ++++++
 .../scheduler/load/LoadTsFileDispatcherImpl.java   |  23 +++
 .../storageengine/dataregion/wal/node/WALNode.java |  16 +-
 .../audit/DataNodeUserDataTransferAuditorTest.java | 136 ++++++++++++++++
 .../execution/exchange/SourceHandleTest.java       |  74 +++++++++
 .../plan/scheduler/AsyncPlanNodeSenderTest.java    |  81 ++++++++++
 .../wal/node/ConsensusReqReaderTest.java           |   2 +
 .../apache/iotdb/commons/i18n/CommonMessages.java  |   4 +
 .../apache/iotdb/commons/i18n/CommonMessages.java  |   3 +
 .../iotdb/commons/audit/AbstractAuditLogger.java   |  49 ++++++
 .../apache/iotdb/commons/audit/AuditEventType.java |   1 +
 .../commons/audit/UserDataTransferAuditEvent.java  |  84 ++++++++++
 ...Type.java => UserDataTransferAuditHandler.java} |  55 +++----
 ...entType.java => UserDataTransferErrorCode.java} |  40 +----
 ....java => UserDataTransferProtectionMethod.java} |  38 +----
 .../commons/audit/AbstractAuditLoggerTest.java     |  19 +++
 .../audit/UserDataTransferAuditEventTest.java      |  57 +++++++
 41 files changed, 1663 insertions(+), 136 deletions(-)

diff --git 
a/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
 
b/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
index 59e3567109b..7f6efd07aeb 100644
--- 
a/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
+++ 
b/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
@@ -333,5 +333,8 @@ public final class IoTConsensusMessages {
   public static final String 
LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = 
"Reserved {} bytes for batch {}-{}, current total usage {}";
   public static final String 
LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}: 
Failed to send idle writer safe-time barrier to {}. status={}";
   public static final String 
LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576
 = "{}: write operation failed. searchIndex: {}. Code: {}, subscriptionQueues: 
{}, this: {}";
+  public static final String
+      
LOG_FAILED_TO_RECORD_A_USER_DATA_TRANSFER_AUDIT_EVENT_CONSENSUS_REPLICATION_WILL_CONTINUE_F215E222
 =
+          "Failed to record a user-data transfer audit event; consensus 
replication will continue.";
 
 }
diff --git 
a/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
 
b/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
index 5164b57b9f5..7cc8ac622d1 100644
--- 
a/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
+++ 
b/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java
@@ -331,5 +331,8 @@ public final class IoTConsensusMessages {
   public static final String 
LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = "预留 {} 
字节给批次 {}-{},当前总使用量 {}";
   public static final String 
LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}:无法向 
{} 发送 idle writer safe-time barrier。状态={}";
   public static final String 
LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576
 = "{}:写入操作失败。searchIndex: {}。Code: {},订阅队列:{},当前对象:{}";
+  public static final String
+      
LOG_FAILED_TO_RECORD_A_USER_DATA_TRANSFER_AUDIT_EVENT_CONSENSUS_REPLICATION_WILL_CONTINUE_F215E222
 =
+          "记录用户数据传送审计事件失败;Consensus 复制将继续。";
 
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
index 834a752be6a..6f67c10d535 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
@@ -49,6 +49,7 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
   private long memorySize = 0;
   private long retainedMemorySize = 0;
   private boolean serializedRequestsBuilt = false;
+  private boolean containsUserData = false;
   private final AtomicLong referenceCnt = new AtomicLong();
 
   public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> 
requests) {
@@ -171,6 +172,15 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
     return this;
   }
 
+  public boolean containsUserData() {
+    return containsUserData;
+  }
+
+  public IndexedConsensusRequest setContainsUserData(boolean containsUserData) 
{
+    this.containsUserData = containsUserData;
+    return this;
+  }
+
   public long getLocalSeq() {
     return searchIndex;
   }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/ConsensusConfig.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/ConsensusConfig.java
index 114a8aed4fc..e3ea222088c 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/ConsensusConfig.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/ConsensusConfig.java
@@ -22,6 +22,7 @@ package org.apache.iotdb.consensus.config;
 import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.commons.audit.TrustedChannelFailureHandler;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
 import org.apache.iotdb.commons.disk.strategy.DirectoryStrategyType;
 
 import java.util.List;
@@ -39,6 +40,8 @@ public class ConsensusConfig {
   private final IoTConsensusV2Config iotConsensusV2Config;
   private final DirectoryStrategyType directoryStrategyType;
   private final TrustedChannelFailureHandler trustedChannelFailureHandler;
+  private final UserDataTransferAuditHandler userDataTransferAuditHandler;
+  private final UserDataTransferAuditClassifier 
userDataTransferAuditClassifier;
 
   private ConsensusConfig(
       TEndPoint thisNode,
@@ -50,7 +53,9 @@ public class ConsensusConfig {
       IoTConsensusConfig iotConsensusConfig,
       IoTConsensusV2Config iotConsensusV2Config,
       DirectoryStrategyType directoryStrategyType,
-      TrustedChannelFailureHandler trustedChannelFailureHandler) {
+      TrustedChannelFailureHandler trustedChannelFailureHandler,
+      UserDataTransferAuditHandler userDataTransferAuditHandler,
+      UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
     this.thisNodeEndPoint = thisNode;
     this.thisNodeId = thisNodeId;
     this.storageDir = storageDir;
@@ -61,6 +66,8 @@ public class ConsensusConfig {
     this.iotConsensusV2Config = iotConsensusV2Config;
     this.directoryStrategyType = directoryStrategyType;
     this.trustedChannelFailureHandler = trustedChannelFailureHandler;
+    this.userDataTransferAuditHandler = userDataTransferAuditHandler;
+    this.userDataTransferAuditClassifier = userDataTransferAuditClassifier;
   }
 
   public TEndPoint getThisNodeEndPoint() {
@@ -103,6 +110,14 @@ public class ConsensusConfig {
     return trustedChannelFailureHandler;
   }
 
+  public UserDataTransferAuditHandler getUserDataTransferAuditHandler() {
+    return userDataTransferAuditHandler;
+  }
+
+  public UserDataTransferAuditClassifier getUserDataTransferAuditClassifier() {
+    return userDataTransferAuditClassifier;
+  }
+
   public static ConsensusConfig.Builder newBuilder() {
     return new ConsensusConfig.Builder();
   }
@@ -121,6 +136,10 @@ public class ConsensusConfig {
         DirectoryStrategyType.MIN_FOLDER_OCCUPIED_SPACE_FIRST_STRATEGY;
     private TrustedChannelFailureHandler trustedChannelFailureHandler =
         TrustedChannelFailureHandler.NO_OP;
+    private UserDataTransferAuditHandler userDataTransferAuditHandler =
+        UserDataTransferAuditHandler.NO_OP;
+    private UserDataTransferAuditClassifier userDataTransferAuditClassifier =
+        UserDataTransferAuditClassifier.NO_USER_DATA;
 
     public ConsensusConfig build() {
       return new ConsensusConfig(
@@ -135,7 +154,9 @@ public class ConsensusConfig {
           Optional.ofNullable(iotConsensusV2Config)
               .orElseGet(() -> IoTConsensusV2Config.newBuilder().build()),
           directoryStrategyType,
-          trustedChannelFailureHandler);
+          trustedChannelFailureHandler,
+          userDataTransferAuditHandler,
+          userDataTransferAuditClassifier);
     }
 
     public Builder setThisNode(TEndPoint thisNode) {
@@ -190,5 +211,21 @@ public class ConsensusConfig {
               .orElse(TrustedChannelFailureHandler.NO_OP);
       return this;
     }
+
+    public Builder setUserDataTransferAuditHandler(
+        UserDataTransferAuditHandler userDataTransferAuditHandler) {
+      this.userDataTransferAuditHandler =
+          Optional.ofNullable(userDataTransferAuditHandler)
+              .orElse(UserDataTransferAuditHandler.NO_OP);
+      return this;
+    }
+
+    public Builder setUserDataTransferAuditClassifier(
+        UserDataTransferAuditClassifier userDataTransferAuditClassifier) {
+      this.userDataTransferAuditClassifier =
+          Optional.ofNullable(userDataTransferAuditClassifier)
+              .orElse(UserDataTransferAuditClassifier.NO_USER_DATA);
+      return this;
+    }
   }
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/UserDataTransferAuditClassifier.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/UserDataTransferAuditClassifier.java
new file mode 100644
index 00000000000..8c8e418cb08
--- /dev/null
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/UserDataTransferAuditClassifier.java
@@ -0,0 +1,47 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.consensus.config;
+
+import org.apache.iotdb.commons.consensus.ConsensusGroupId;
+import org.apache.iotdb.commons.request.IConsensusRequest;
+
+@FunctionalInterface
+public interface UserDataTransferAuditClassifier {
+
+  UserDataTransferAuditClassifier NO_USER_DATA =
+      new UserDataTransferAuditClassifier() {
+        @Override
+        public boolean containsUserData(ConsensusGroupId groupId, 
IConsensusRequest request) {
+          return false;
+        }
+
+        @Override
+        public boolean containsUserData(ConsensusGroupId groupId) {
+          return false;
+        }
+      };
+
+  boolean containsUserData(ConsensusGroupId groupId, IConsensusRequest 
request);
+
+  /** Classifies a whole consensus group when the transfer has no individual 
request to inspect. */
+  default boolean containsUserData(ConsensusGroupId groupId) {
+    return true;
+  }
+}
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
index 477d8a5cb11..ede52c92ed3 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.consensus.iot;
 
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.commons.concurrent.ThreadName;
@@ -45,6 +46,7 @@ import org.apache.iotdb.consensus.common.DataSet;
 import org.apache.iotdb.consensus.common.Peer;
 import org.apache.iotdb.consensus.config.ConsensusConfig;
 import org.apache.iotdb.consensus.config.IoTConsensusConfig;
+import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
 import org.apache.iotdb.consensus.exception.ConsensusException;
 import 
org.apache.iotdb.consensus.exception.ConsensusGroupAlreadyExistException;
 import org.apache.iotdb.consensus.exception.ConsensusGroupModifyPeerException;
@@ -104,6 +106,8 @@ public class IoTConsensus implements IConsensus {
       new ConcurrentHashMap<>();
   private final IoTConsensusRPCService service;
   private final RegisterManager registerManager = new RegisterManager();
+  private final UserDataTransferAuditHandler userDataTransferAuditHandler;
+  private final UserDataTransferAuditClassifier 
userDataTransferAuditClassifier;
   private volatile IoTConsensusConfig config;
 
   /**
@@ -131,6 +135,8 @@ public class IoTConsensus implements IConsensus {
     this.recvSnapshotDirs = config.getRecvSnapshotDirs();
     this.recvFolderStrategyType = config.getDirectoryStrategyType();
     this.config = config.getIotConsensusConfig();
+    this.userDataTransferAuditHandler = 
config.getUserDataTransferAuditHandler();
+    this.userDataTransferAuditClassifier = 
config.getUserDataTransferAuditClassifier();
     this.registry = registry;
     this.service =
         new IoTConsensusRPCService(
@@ -207,7 +213,9 @@ public class IoTConsensus implements IConsensus {
                   backgroundTaskService,
                   clientManager,
                   syncClientManager,
-                  config);
+                  config,
+                  userDataTransferAuditHandler,
+                  userDataTransferAuditClassifier);
           stateMachineMap.put(consensusGroupId, consensus);
         }
       } catch (DiskSpaceInsufficientException e) {
@@ -322,7 +330,9 @@ public class IoTConsensus implements IConsensus {
                             backgroundTaskService,
                             clientManager,
                             syncClientManager,
-                            config);
+                            config,
+                            userDataTransferAuditHandler,
+                            userDataTransferAuditClassifier);
                   } catch (DiskSpaceInsufficientException e) {
                     throw new RuntimeException(e);
                   }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
index e074e7204ee..6e86ec72a3e 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
@@ -21,6 +21,9 @@ package org.apache.iotdb.consensus.iot;
 
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
+import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.client.exception.ClientManagerException;
 import org.apache.iotdb.commons.consensus.ConsensusGroupId;
@@ -42,6 +45,7 @@ import org.apache.iotdb.consensus.common.Peer;
 import 
org.apache.iotdb.consensus.common.request.DeserializedBatchIndexedConsensusRequest;
 import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
 import org.apache.iotdb.consensus.config.IoTConsensusConfig;
+import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
 import org.apache.iotdb.consensus.exception.ConsensusGroupModifyPeerException;
 import org.apache.iotdb.consensus.i18n.ConsensusMessages;
 import org.apache.iotdb.consensus.i18n.IoTConsensusMessages;
@@ -158,6 +162,8 @@ public class IoTConsensusServerImpl {
   private final ScheduledExecutorService backgroundTaskService;
   private final IoTConsensusRateLimiter ioTConsensusRateLimiter =
       IoTConsensusRateLimiter.getInstance();
+  private final UserDataTransferAuditHandler userDataTransferAuditHandler;
+  private final UserDataTransferAuditClassifier 
userDataTransferAuditClassifier;
   private IndexedConsensusRequest lastConsensusRequest;
 
   // Subscription queues receive IndexedConsensusRequest in real-time from 
write(),
@@ -196,6 +202,35 @@ public class IoTConsensusServerImpl {
       IClientManager<TEndPoint, SyncIoTConsensusServiceClient> 
syncClientManager,
       IoTConsensusConfig config)
       throws DiskSpaceInsufficientException {
+    this(
+        storageDir,
+        recvSnapshotDirs,
+        recvFolderStrategyType,
+        thisNode,
+        configuration,
+        stateMachine,
+        backgroundTaskService,
+        clientManager,
+        syncClientManager,
+        config,
+        UserDataTransferAuditHandler.NO_OP,
+        UserDataTransferAuditClassifier.NO_USER_DATA);
+  }
+
+  public IoTConsensusServerImpl(
+      String storageDir,
+      List<String> recvSnapshotDirs,
+      DirectoryStrategyType recvFolderStrategyType,
+      Peer thisNode,
+      Collection<Peer> configuration,
+      IStateMachine stateMachine,
+      ScheduledExecutorService backgroundTaskService,
+      IClientManager<TEndPoint, AsyncIoTConsensusServiceClient> clientManager,
+      IClientManager<TEndPoint, SyncIoTConsensusServiceClient> 
syncClientManager,
+      IoTConsensusConfig config,
+      UserDataTransferAuditHandler userDataTransferAuditHandler,
+      UserDataTransferAuditClassifier userDataTransferAuditClassifier)
+      throws DiskSpaceInsufficientException {
     this.active = true;
     this.storageDir = storageDir;
     List<String> snapshotDirs = new ArrayList<>();
@@ -215,6 +250,8 @@ public class IoTConsensusServerImpl {
     this.configuration.addAll(configuration);
     this.backgroundTaskService = backgroundTaskService;
     this.config = config;
+    this.userDataTransferAuditHandler = userDataTransferAuditHandler;
+    this.userDataTransferAuditClassifier = userDataTransferAuditClassifier;
     this.consensusGroupId = thisNode.getGroupId().toString();
     this.consensusReqReader =
         (ConsensusReqReader) stateMachine.read(new 
GetConsensusReqReaderPlan());
@@ -391,6 +428,9 @@ public class IoTConsensusServerImpl {
   public void transmitSnapshot(Peer targetPeer) throws 
ConsensusGroupModifyPeerException {
     File snapshotDir = new File(storageDir, newSnapshotDirName);
     List<File> snapshotPaths = stateMachine.getSnapshotFiles(snapshotDir);
+    final boolean auditSnapshotTransfer =
+        shouldAuditSnapshotTransfer(
+            userDataTransferAuditHandler, userDataTransferAuditClassifier, 
thisNode.getGroupId());
     long snapshotSizeSum = 0;
     for (File file : snapshotPaths) {
       snapshotSizeSum += file.length();
@@ -433,7 +473,19 @@ public class IoTConsensusServerImpl {
             TSendSnapshotFragmentReq req = 
reader.next().toTSendSnapshotFragmentReq();
             
req.setConsensusGroupId(targetPeer.getGroupId().convertToTConsensusGroupId());
             
ioTConsensusRateLimiter.acquireTransitDataSizeWithRateLimiter(req.getChunkLength());
-            TSendSnapshotFragmentRes res = client.sendSnapshotFragment(req);
+            final TSendSnapshotFragmentRes res;
+            try {
+              res = client.sendSnapshotFragment(req);
+              recordSnapshotTransferAttempt(
+                  targetPeer,
+                  auditSnapshotTransfer,
+                  isSuccess(res.getStatus()),
+                  isSuccess(res.getStatus()) ? null : 
String.valueOf(res.getStatus().getCode()),
+                  null);
+            } catch (Exception e) {
+              recordSnapshotTransferAttempt(targetPeer, auditSnapshotTransfer, 
false, null, e);
+              throw e;
+            }
             if (!isSuccess(res.getStatus())) {
               throw new ConsensusGroupModifyPeerException(
                   
String.format(IoTConsensusMessages.SNAPSHOT_TRANSMISSION_ERROR, targetPeer));
@@ -470,6 +522,41 @@ public class IoTConsensusServerImpl {
         snapshotDir);
   }
 
+  private void recordSnapshotTransferAttempt(
+      Peer targetPeer,
+      boolean auditSnapshotTransfer,
+      boolean success,
+      String errorCode,
+      Throwable error) {
+    if (!auditSnapshotTransfer) {
+      return;
+    }
+    try {
+      userDataTransferAuditHandler.onAttempt(
+          new UserDataTransferAuditEvent(
+              thisNode.getEndpoint(),
+              thisNode.getEndpoint(),
+              targetPeer.getEndpoint(),
+              
UserDataTransferProtectionMethod.fromTlsEnabled(config.getRpc().isEnableSSL()),
+              success,
+              errorCode != null ? errorCode : error == null ? null : 
error.getClass().getName()));
+    } catch (RuntimeException ignored) {
+      // Audit recording must not affect snapshot transmission.
+    }
+  }
+
+  static boolean shouldAuditSnapshotTransfer(
+      UserDataTransferAuditHandler auditHandler,
+      UserDataTransferAuditClassifier auditClassifier,
+      ConsensusGroupId groupId) {
+    try {
+      return auditHandler.isEnabled() && 
auditClassifier.containsUserData(groupId);
+    } catch (RuntimeException ignored) {
+      // Audit classification must not affect snapshot transmission.
+      return false;
+    }
+  }
+
   public void receiveSnapshotFragment(
       String snapshotId, String originalFilePath, ByteBuffer fileChunk, long 
fileOffset)
       throws ConsensusGroupModifyPeerException {
@@ -940,7 +1027,9 @@ public class IoTConsensusServerImpl {
           new IoTProgressIndex(thisNode.getNodeId(), searchIndex.get() + 1);
       ((ComparableConsensusRequest) 
request).setProgressIndex(iotProgressIndex);
     }
-    return new IndexedConsensusRequest(searchIndex.get() + 1, 
Collections.singletonList(request))
+    final List<IConsensusRequest> requests = 
Collections.singletonList(request);
+    return new IndexedConsensusRequest(searchIndex.get() + 1, requests)
+        .setContainsUserData(containsUserData(request))
         .setPhysicalTime(assignPhysicalTimeInMs())
         .setNodeId(thisNode.getNodeId());
   }
@@ -960,6 +1049,33 @@ public class IoTConsensusServerImpl {
     return req;
   }
 
+  public boolean containsUserData() {
+    try {
+      return 
userDataTransferAuditClassifier.containsUserData(thisNode.getGroupId());
+    } catch (RuntimeException ignored) {
+      // Classification is advisory and must not affect consensus replication.
+      return false;
+    }
+  }
+
+  public boolean containsUserData(List<IConsensusRequest> requests) {
+    for (IConsensusRequest request : requests) {
+      if (containsUserData(request)) {
+        return true;
+      }
+    }
+    return false;
+  }
+
+  private boolean containsUserData(IConsensusRequest request) {
+    try {
+      return 
userDataTransferAuditClassifier.containsUserData(thisNode.getGroupId(), 
request);
+    } catch (RuntimeException ignored) {
+      // Classification is advisory and must not affect consensus replication.
+      return false;
+    }
+  }
+
   public TSStatus syncIdleWriterSafeTimeBarrierToPeer(final Peer targetPeer) {
     final long safePhysicalTime = assignPhysicalTimeInMs();
     final long safeLocalSeq = searchIndex.get();
@@ -1121,6 +1237,10 @@ public class IoTConsensusServerImpl {
     return thisNode;
   }
 
+  public UserDataTransferAuditHandler getUserDataTransferAuditHandler() {
+    return userDataTransferAuditHandler;
+  }
+
   public List<Peer> getConfiguration() {
     List<Peer> result = new ArrayList<>(configuration);
     Collections.sort(result);
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandler.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandler.java
index 64034144b93..615ecd6b07c 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandler.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandler.java
@@ -19,7 +19,11 @@
 
 package org.apache.iotdb.consensus.iot.client;
 
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
+import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
 import org.apache.iotdb.commons.utils.RetryUtils;
 import org.apache.iotdb.consensus.i18n.IoTConsensusMessages;
 import org.apache.iotdb.consensus.iot.logdispatcher.Batch;
@@ -41,7 +45,7 @@ import java.util.stream.Collectors;
 
 public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRes> {
 
-  private final Logger logger = 
LoggerFactory.getLogger(DispatchLogHandler.class);
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(DispatchLogHandler.class);
 
   private final LogDispatcherThread thread;
   private final Batch batch;
@@ -63,6 +67,7 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
 
   @Override
   public void onComplete(TSyncLogEntriesRes response) {
+    recordTransferAttempt(response);
     if (response.getStatuses().stream()
         .anyMatch(status -> RetryUtils.needRetryForWrite(status.getCode()))) {
       List<String> retryStatusMessages =
@@ -73,14 +78,14 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
 
       String messages = String.join(", ", retryStatusMessages);
       if (++retryCount == 1) {
-        logger.warn(
+        LOGGER.warn(
             IoTConsensusMessages.CANNOT_SEND_TO_PEER,
             batch,
             thread.getPeer(),
             retryCount,
             messages);
       } else {
-        logger.debug(
+        LOGGER.debug(
             IoTConsensusMessages.CANNOT_SEND_TO_PEER,
             batch,
             thread.getPeer(),
@@ -89,13 +94,13 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
       }
       sleepCorrespondingTimeAndRetryAsynchronous();
     } else {
-      if (logger.isDebugEnabled()) {
+      if (LOGGER.isDebugEnabled()) {
         boolean containsError =
             response.getStatuses().stream()
                 .anyMatch(
                     status -> status.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode());
         if (containsError) {
-          logger.debug(
+          LOGGER.debug(
               IoTConsensusMessages.SEND_COMPLETE_BUT_CONTAINS_ERROR,
               batch,
               thread.getPeer(),
@@ -113,18 +118,19 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
 
   @Override
   public void onError(Exception exception) {
+    recordTransferAttempt(false, null, exception);
     ++retryCount;
     Throwable rootCause = ExceptionUtils.getRootCause(exception);
     final Throwable actualCause = rootCause == null ? exception : rootCause;
     if (retryCount == 1) {
-      logger.warn(
+      LOGGER.warn(
           IoTConsensusMessages.CANNOT_SEND_TO_PEER_ON_ERROR,
           batch,
           thread.getPeer(),
           retryCount,
           actualCause.toString());
     } else {
-      logger.debug(
+      LOGGER.debug(
           IoTConsensusMessages.CANNOT_SEND_TO_PEER_ON_ERROR,
           batch,
           thread.getPeer(),
@@ -134,7 +140,7 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
     // skip TApplicationException caused by follower
     if (actualCause instanceof TApplicationException) {
       completeBatch(batch);
-      logger.warn(IoTConsensusMessages.SKIP_RETRY_TAPPLICATION_EXCEPTION, 
batch);
+      LOGGER.warn(IoTConsensusMessages.SKIP_RETRY_TAPPLICATION_EXCEPTION, 
batch);
       logDispatcherThreadMetrics.recordSyncLogTimePerRequest(System.nanoTime() 
- createTime);
       return;
     }
@@ -153,7 +159,7 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
         .schedule(
             () -> {
               if (thread.isStopped()) {
-                logger.debug(
+                LOGGER.debug(
                     IoTConsensusMessages.LOG_DISPATCHER_STOPPED_NO_RETRY,
                     thread.getPeer(),
                     batch,
@@ -172,4 +178,111 @@ public class DispatchLogHandler implements 
AsyncMethodCallback<TSyncLogEntriesRe
     // removeBatch
     thread.updateSafelyDeletedSearchIndex();
   }
+
+  private void recordTransferAttempt(boolean success, String errorCode, 
Throwable error) {
+    try {
+      final UserDataTransferAuditHandler auditHandler =
+          thread.getImpl().getUserDataTransferAuditHandler();
+      if (!batch.containsUserData() || !auditHandler.isEnabled()) {
+        return;
+      }
+      recordTransferAttempt(
+          auditHandler,
+          batch,
+          thread.getImpl().getThisNode().getEndpoint(),
+          thread.getPeer().getEndpoint(),
+          UserDataTransferProtectionMethod.fromTlsEnabled(
+              thread.getConfig().getRpc().isEnableSSL()),
+          success,
+          errorCode,
+          error);
+    } catch (RuntimeException auditFailure) {
+      warnAuditFailure(auditFailure);
+    }
+  }
+
+  private void recordTransferAttempt(TSyncLogEntriesRes response) {
+    try {
+      final UserDataTransferAuditHandler auditHandler =
+          thread.getImpl().getUserDataTransferAuditHandler();
+      if (!batch.containsUserData() || !auditHandler.isEnabled()) {
+        return;
+      }
+      recordTransferAttempt(
+          auditHandler,
+          batch,
+          thread.getImpl().getThisNode().getEndpoint(),
+          thread.getPeer().getEndpoint(),
+          UserDataTransferProtectionMethod.fromTlsEnabled(
+              thread.getConfig().getRpc().isEnableSSL()),
+          response);
+    } catch (RuntimeException auditFailure) {
+      warnAuditFailure(auditFailure);
+    }
+  }
+
+  static void recordTransferAttempt(
+      UserDataTransferAuditHandler auditHandler,
+      Batch batch,
+      TEndPoint source,
+      TEndPoint target,
+      UserDataTransferProtectionMethod protectionMethod,
+      TSyncLogEntriesRes response) {
+    try {
+      if (!batch.containsUserData() || !auditHandler.isEnabled()) {
+        return;
+      }
+      // One batch RPC is one physical transfer attempt, so keep one 
representative error value in
+      // the minimum audit record instead of concatenating an unbounded number 
of response details.
+      final TSStatus firstFailedStatus =
+          response.getStatuses().stream()
+              .filter(status -> status.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode())
+              .findFirst()
+              .orElse(null);
+      recordTransferAttempt(
+          auditHandler,
+          batch,
+          source,
+          target,
+          protectionMethod,
+          firstFailedStatus == null,
+          firstFailedStatus == null ? null : 
String.valueOf(firstFailedStatus.getCode()),
+          null);
+    } catch (RuntimeException auditFailure) {
+      warnAuditFailure(auditFailure);
+    }
+  }
+
+  static void recordTransferAttempt(
+      UserDataTransferAuditHandler auditHandler,
+      Batch batch,
+      TEndPoint source,
+      TEndPoint target,
+      UserDataTransferProtectionMethod protectionMethod,
+      boolean success,
+      String errorCode,
+      Throwable error) {
+    try {
+      if (!batch.containsUserData() || !auditHandler.isEnabled()) {
+        return;
+      }
+      auditHandler.onAttempt(
+          new UserDataTransferAuditEvent(
+              source,
+              source,
+              target,
+              protectionMethod,
+              success,
+              errorCode != null ? errorCode : error == null ? null : 
error.getClass().getName()));
+    } catch (RuntimeException auditFailure) {
+      warnAuditFailure(auditFailure);
+    }
+  }
+
+  private static void warnAuditFailure(RuntimeException auditFailure) {
+    LOGGER.warn(
+        IoTConsensusMessages
+            
.LOG_FAILED_TO_RECORD_A_USER_DATA_TRANSFER_AUDIT_EVENT_CONSENSUS_REPLICATION_WILL_CONTINUE_F215E222,
+        auditFailure);
+  }
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java
index 72b68ab96ac..55569b8a34f 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java
@@ -39,6 +39,7 @@ public class Batch {
   private long memorySize;
   // indicates whether this batch has been successfully synchronized to 
another node
   private boolean synced;
+  private boolean containsUserData;
 
   public Batch(IoTConsensusConfig config) {
     this.config = config;
@@ -55,11 +56,16 @@ public class Batch {
   }
 
   public void addTLogEntry(TLogEntry entry) {
+    addTLogEntry(entry, false);
+  }
+
+  public void addTLogEntry(TLogEntry entry, boolean containsUserData) {
     logEntries.add(entry);
     if (entry.fromWAL) {
       logEntriesNumFromWAL++;
     }
     memorySize += entry.getMemorySize();
+    this.containsUserData |= containsUserData;
   }
 
   public boolean canAccumulate() {
@@ -107,6 +113,10 @@ public class Batch {
     return logEntriesNumFromWAL;
   }
 
+  public boolean containsUserData() {
+    return containsUserData;
+  }
+
   @Override
   public String toString() {
     return "Batch{"
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
index cd9d7eea49f..1ff6579e0de 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
@@ -651,6 +651,7 @@ public class LogDispatcher {
           currentIndex,
           maxIndex);
       boolean hasCorruptedData = false;
+      final boolean consensusGroupContainsUserData = impl.containsUserData();
       // targetIndex is the index of request that we need to find
       long targetIndex = currentIndex;
       // Even if there is no WAL files, these code won't produce error.
@@ -675,6 +676,9 @@ public class LogDispatcher {
           hasCorruptedData = true;
         }
         targetIndex = data.getSearchIndex() + 1;
+        // The WAL reader derives this bit directly from the entry type. Apply 
the group-level
+        // exclusion here without deserializing the request solely for audit 
classification.
+        data.setContainsUserData(consensusGroupContainsUserData && 
data.containsUserData());
         data.buildSerializedRequests();
         // construct request from wal
         TLogEntry logEntry =
@@ -682,7 +686,7 @@ public class LogDispatcher {
                 data.getSerializedRequests(), data.getSearchIndex(), true, 
data.getMemorySize());
         logEntry.setRoutingEpoch(data.getRoutingEpoch());
         logEntry.setPhysicalTime(data.getPhysicalTime());
-        logBatches.addTLogEntry(logEntry);
+        logBatches.addTLogEntry(logEntry, data.containsUserData());
       }
       // In the case of corrupt Data, we return true so that we can send a 
batch as soon as
       // possible, avoiding potential duplication
@@ -699,7 +703,7 @@ public class LogDispatcher {
               request.getMemorySize());
       logEntry.setRoutingEpoch(request.getRoutingEpoch());
       logEntry.setPhysicalTime(request.getPhysicalTime());
-      logBatches.addTLogEntry(logEntry);
+      logBatches.addTLogEntry(logEntry, request.containsUserData());
     }
   }
 
diff --git 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImplTest.java
 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImplTest.java
index 24f8e8dc6bd..cc33cf1573a 100644
--- 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImplTest.java
+++ 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImplTest.java
@@ -20,10 +20,14 @@
 package org.apache.iotdb.consensus.iot;
 
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
+import org.apache.iotdb.commons.consensus.ConsensusGroupId;
 import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.disk.strategy.DirectoryStrategyType;
+import org.apache.iotdb.commons.request.IConsensusRequest;
 import org.apache.iotdb.consensus.common.Peer;
 import org.apache.iotdb.consensus.config.IoTConsensusConfig;
+import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
 import org.apache.iotdb.consensus.iot.util.TestStateMachine;
 
 import org.junit.Rule;
@@ -44,6 +48,7 @@ import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertTrue;
 
 public class IoTConsensusServerImplTest {
@@ -54,6 +59,37 @@ public class IoTConsensusServerImplTest {
 
   @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
 
+  @Test
+  public void testSnapshotTransferAppliesGroupClassifier() {
+    final DataRegionId groupId = new DataRegionId(1);
+    final UserDataTransferAuditHandler enabledHandler = event -> {};
+    final UserDataTransferAuditClassifier auditDatabaseClassifier =
+        new UserDataTransferAuditClassifier() {
+          @Override
+          public boolean containsUserData(
+              ConsensusGroupId ignoredGroupId, IConsensusRequest request) {
+            return false;
+          }
+
+          @Override
+          public boolean containsUserData(ConsensusGroupId ignoredGroupId) {
+            return false;
+          }
+        };
+
+    assertTrue(
+        IoTConsensusServerImpl.shouldAuditSnapshotTransfer(
+            enabledHandler, (ignoredGroupId, request) -> false, groupId));
+    assertFalse(
+        IoTConsensusServerImpl.shouldAuditSnapshotTransfer(
+            enabledHandler, auditDatabaseClassifier, groupId));
+    assertFalse(
+        IoTConsensusServerImpl.shouldAuditSnapshotTransfer(
+            UserDataTransferAuditHandler.NO_OP,
+            UserDataTransferAuditClassifier.NO_USER_DATA,
+            groupId));
+  }
+
   /**
    * Verifies that configuration snapshots can be read while several writers 
concurrently add and
    * remove distinct peers. Every snapshot must remain duplicate-free and 
sorted, and the final
diff --git 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandlerTest.java
 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandlerTest.java
new file mode 100644
index 00000000000..13736903902
--- /dev/null
+++ 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/client/DispatchLogHandlerTest.java
@@ -0,0 +1,173 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.consensus.iot.client;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
+import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
+import org.apache.iotdb.consensus.config.IoTConsensusConfig;
+import org.apache.iotdb.consensus.iot.logdispatcher.Batch;
+import org.apache.iotdb.consensus.iot.thrift.TLogEntry;
+import org.apache.iotdb.consensus.iot.thrift.TSyncLogEntriesRes;
+
+import org.junit.Test;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public class DispatchLogHandlerTest {
+
+  private static final TEndPoint SOURCE = new TEndPoint("127.0.0.1", 10740);
+  private static final TEndPoint TARGET = new TEndPoint("127.0.0.2", 10740);
+
+  @Test
+  public void testRecordsSuccessRemoteFailureAndException() {
+    final List<UserDataTransferAuditEvent> events = new ArrayList<>();
+    final Batch batch = createBatch(true);
+
+    DispatchLogHandler.recordTransferAttempt(
+        events::add, batch, SOURCE, TARGET, 
UserDataTransferProtectionMethod.TLS, true, null, null);
+    DispatchLogHandler.recordTransferAttempt(
+        events::add,
+        batch,
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.TLS,
+        false,
+        "500",
+        null);
+    DispatchLogHandler.recordTransferAttempt(
+        events::add,
+        batch,
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.TLS,
+        false,
+        null,
+        new IOException());
+
+    assertEquals(3, events.size());
+    assertTrue(events.get(0).isSuccess());
+    assertEquals(SOURCE, events.get(0).getInitiator());
+    assertEquals(SOURCE, events.get(0).getSource());
+    assertEquals(TARGET, events.get(0).getTarget());
+    assertEquals(UserDataTransferProtectionMethod.TLS, 
events.get(0).getProtectionMethod());
+    assertFalse(events.get(1).isSuccess());
+    assertEquals("500", events.get(1).getError());
+    assertEquals(IOException.class.getName(), events.get(2).getError());
+  }
+
+  @Test
+  public void testSkipsBatchWithoutUserData() {
+    final List<UserDataTransferAuditEvent> events = new ArrayList<>();
+
+    DispatchLogHandler.recordTransferAttempt(
+        events::add,
+        createBatch(false),
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.NONE,
+        true,
+        null,
+        null);
+
+    assertTrue(events.isEmpty());
+  }
+
+  @Test
+  public void testSkipDoesNotInspectResponseStatuses() {
+    final AtomicInteger statusAccessCount = new AtomicInteger();
+    final TSyncLogEntriesRes response =
+        new TSyncLogEntriesRes() {
+          @Override
+          public List<TSStatus> getStatuses() {
+            statusAccessCount.incrementAndGet();
+            return super.getStatuses();
+          }
+        };
+
+    DispatchLogHandler.recordTransferAttempt(
+        event -> {},
+        createBatch(false),
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.NONE,
+        response);
+    DispatchLogHandler.recordTransferAttempt(
+        UserDataTransferAuditHandler.NO_OP,
+        createBatch(true),
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.NONE,
+        response);
+
+    assertEquals(0, statusAccessCount.get());
+  }
+
+  @Test
+  public void testAuditHandlerFailureDoesNotEscape() {
+    DispatchLogHandler.recordTransferAttempt(
+        event -> {
+          throw new IllegalStateException();
+        },
+        createBatch(true),
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.NONE,
+        true,
+        null,
+        null);
+
+    DispatchLogHandler.recordTransferAttempt(
+        new UserDataTransferAuditHandler() {
+          @Override
+          public void onAttempt(UserDataTransferAuditEvent event) {
+            // Do nothing.
+          }
+
+          @Override
+          public boolean isEnabled() {
+            throw new IllegalStateException();
+          }
+        },
+        createBatch(true),
+        SOURCE,
+        TARGET,
+        UserDataTransferProtectionMethod.NONE,
+        true,
+        null,
+        null);
+  }
+
+  private static Batch createBatch(boolean containsUserData) {
+    final Batch batch = new Batch(IoTConsensusConfig.newBuilder().build());
+    batch.addTLogEntry(new TLogEntry().setSearchIndex(1).setMemorySize(1), 
containsUserData);
+    batch.buildIndex();
+    return batch;
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index 3feefd72a6f..3d129ba8f17 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -1398,6 +1398,8 @@ public final class DataNodeQueryMessages {
       "{} failed to pull TsBlocks [{}] to [{}] from SinkHandle {}, channel 
index {},";
   public static final String FAILED_TO_GET_DATA_BLOCK =
       "failed to get data block [{}, {}), attempt times: {}";
+  public static final String 
EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33 =
+      "Unexpected data block response size.";
   public static final String FAILED_TO_SEND_ACK_DATA_BLOCK_EVENT =
       "failed to send ack data block event [{}, {}), attempt times: {}";
   public static final String SEND_CLOSE_SINK_CHANNEL_EVENT_FAILED =
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index 0d4e7cc0071..88b7b0edb90 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -1379,6 +1379,8 @@ public final class DataNodeQueryMessages {
       "{} 从 SinkHandle {} 的通道索引 {} 拉取 TsBlocks [{}] 到 [{}] 失败,";
   public static final String FAILED_TO_GET_DATA_BLOCK =
       "获取数据块 [{}, {}) 失败,尝试次数:{}";
+  public static final String 
EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33 =
+      "数据块响应数量异常。";
   public static final String FAILED_TO_SEND_ACK_DATA_BLOCK_EVENT =
       "发送数据块确认事件 [{}, {}) 失败,尝试次数:{}";
   public static final String SEND_CLOSE_SINK_CHANNEL_EVENT_FAILED =
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditor.java
new file mode 100644
index 00000000000..0ee38058066
--- /dev/null
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditor.java
@@ -0,0 +1,135 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.audit;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditEvent;
+import org.apache.iotdb.commons.audit.UserDataTransferProtectionMethod;
+import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.commons.consensus.ConsensusGroupId;
+import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.commons.request.IConsensusRequest;
+import org.apache.iotdb.commons.schema.table.Audit;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.ObjectNode;
+import org.apache.iotdb.db.storageengine.StorageEngine;
+import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
+
+import javax.annotation.Nullable;
+
+public final class DataNodeUserDataTransferAuditor {
+
+  private static final CommonConfig COMMON_CONFIG = 
CommonDescriptor.getInstance().getConfig();
+
+  private DataNodeUserDataTransferAuditor() {}
+
+  public static boolean isEnabled() {
+    return COMMON_CONFIG.isEnableAuditLog();
+  }
+
+  public static boolean isEnabledFor(ConsensusGroupId consensusGroupId) {
+    try {
+      return isEnabled() && containsUserData(consensusGroupId);
+    } catch (RuntimeException ignored) {
+      // Classification is advisory and must not affect user data transfer.
+      return false;
+    }
+  }
+
+  public static void record(
+      TEndPoint initiator,
+      TEndPoint source,
+      TEndPoint target,
+      boolean success,
+      @Nullable String errorCode,
+      @Nullable Throwable error) {
+    try {
+      if (!isEnabled()) {
+        return;
+      }
+      DNAuditLogger.getInstance()
+          .recordUserDataTransferAuditLog(
+              new UserDataTransferAuditEvent(
+                  initiator,
+                  source,
+                  target,
+                  UserDataTransferProtectionMethod.fromTlsEnabled(
+                      COMMON_CONFIG.isEnableInternalSSL()),
+                  success,
+                  errorCode != null
+                      ? errorCode
+                      : error == null ? null : error.getClass().getName()));
+    } catch (RuntimeException ignored) {
+      // Audit recording must not affect user data transfer.
+    }
+  }
+
+  public static boolean containsUserData(
+      ConsensusGroupId consensusGroupId, IConsensusRequest request) {
+    if (!(consensusGroupId instanceof DataRegionId)) {
+      return false;
+    }
+    final DataRegion dataRegion =
+        StorageEngine.getInstance().getDataRegion((DataRegionId) 
consensusGroupId);
+    return dataRegion != null && 
containsUserData(dataRegion.getDatabaseName(), request);
+  }
+
+  public static boolean containsUserData(ConsensusGroupId consensusGroupId) {
+    if (!(consensusGroupId instanceof DataRegionId)) {
+      return false;
+    }
+    final DataRegion dataRegion =
+        StorageEngine.getInstance().getDataRegion((DataRegionId) 
consensusGroupId);
+    return dataRegion != null && 
containsUserData(dataRegion.getDatabaseName());
+  }
+
+  static boolean containsUserData(String database) {
+    return !Audit.isAuditDatabase(database);
+  }
+
+  static boolean containsUserData(String database, IConsensusRequest request) {
+    if (!containsUserData(database)) {
+      return false;
+    }
+    try {
+      return request instanceof PlanNode && containsUserData((PlanNode) 
request);
+    } catch (RuntimeException ignored) {
+      // Classification is advisory and must not affect consensus replication.
+      return false;
+    }
+  }
+
+  public static boolean containsUserData(PlanNode node) {
+    if (node instanceof InsertNode || node instanceof ObjectNode) {
+      return true;
+    }
+    if (node.getChildren() == null) {
+      return false;
+    }
+    for (PlanNode child : node.getChildren()) {
+      if (containsUserData(child)) {
+        return true;
+      }
+    }
+    return false;
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
index fbda403b320..ba08b7379bc 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.consensus;
 
 import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.commons.audit.UserDataTransferAuditHandler;
 import org.apache.iotdb.commons.conf.CommonConfig;
 import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.consensus.ConsensusGroupId;
@@ -29,6 +30,7 @@ import 
org.apache.iotdb.commons.disk.strategy.DirectoryStrategyType;
 import org.apache.iotdb.commons.memory.IMemoryBlock;
 import org.apache.iotdb.commons.memory.MemoryBlockType;
 import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
+import org.apache.iotdb.commons.request.IConsensusRequest;
 import org.apache.iotdb.commons.utils.TestOnly;
 import org.apache.iotdb.consensus.ConsensusFactory;
 import org.apache.iotdb.consensus.IConsensus;
@@ -39,7 +41,9 @@ import org.apache.iotdb.consensus.config.IoTConsensusV2Config;
 import org.apache.iotdb.consensus.config.IoTConsensusV2Config.ReplicateMode;
 import org.apache.iotdb.consensus.config.RatisConfig;
 import org.apache.iotdb.consensus.config.RatisConfig.Snapshot;
+import org.apache.iotdb.consensus.config.UserDataTransferAuditClassifier;
 import org.apache.iotdb.db.audit.DNAuditLogger;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.db.conf.DataNodeMemoryConfig;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -67,6 +71,18 @@ import java.util.concurrent.TimeUnit;
 public class DataRegionConsensusImpl {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(DataRegionConsensusImpl.class);
+  private static final UserDataTransferAuditClassifier 
USER_DATA_TRANSFER_AUDIT_CLASSIFIER =
+      new UserDataTransferAuditClassifier() {
+        @Override
+        public boolean containsUserData(ConsensusGroupId groupId, 
IConsensusRequest request) {
+          return DataNodeUserDataTransferAuditor.containsUserData(groupId, 
request);
+        }
+
+        @Override
+        public boolean containsUserData(ConsensusGroupId groupId) {
+          return DataNodeUserDataTransferAuditor.containsUserData(groupId);
+        }
+      };
 
   private DataRegionConsensusImpl() {
     // do nothing
@@ -143,6 +159,14 @@ public class DataRegionConsensusImpl {
           .setThisNode(new TEndPoint(CONF.getInternalAddress(), 
CONF.getDataRegionConsensusPort()))
           .setTrustedChannelFailureHandler(
               
DNAuditLogger.getInstance()::recordTrustedChannelFailureAuditLogIfNecessary)
+          .setUserDataTransferAuditHandler(
+              COMMON_CONF.isEnableAuditLog()
+                  ? DNAuditLogger.getInstance()::recordUserDataTransferAuditLog
+                  : UserDataTransferAuditHandler.NO_OP)
+          .setUserDataTransferAuditClassifier(
+              COMMON_CONF.isEnableAuditLog()
+                  ? USER_DATA_TRANSFER_AUDIT_CLASSIFIER
+                  : UserDataTransferAuditClassifier.NO_USER_DATA)
           .setStorageDir(CONF.getDataRegionConsensusDir())
           .setRecvSnapshotDirs(Arrays.asList(CONF.getLocalDataDirs()))
           // IoTConsensus always balances received snapshot files by least 
occupied space,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2AsyncSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2AsyncSink.java
index 05e14ed08ae..327309005d8 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2AsyncSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2AsyncSink.java
@@ -22,9 +22,11 @@ package 
org.apache.iotdb.db.pipe.sink.protocol.iotconsensusv2;
 import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
 import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.client.async.AsyncIoTConsensusV2ServiceClient;
 import org.apache.iotdb.commons.consensus.ConsensusGroupId;
+import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import 
org.apache.iotdb.commons.consensus.iotv2.container.IoTV2GlobalComponentContainer;
 import 
org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkRetryTimesConfigurableException;
@@ -37,6 +39,7 @@ import 
org.apache.iotdb.consensus.iotconsensusv2.thrift.TIoTConsensusV2TransferR
 import org.apache.iotdb.consensus.pipe.consensuspipe.ConsensusPipeName;
 import org.apache.iotdb.consensus.pipe.consensuspipe.ConsensusPipeSink;
 import org.apache.iotdb.consensus.pipe.metric.IoTConsensusV2SyncLagManager;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodePipeMessages;
@@ -63,6 +66,7 @@ import 
org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
 import org.apache.iotdb.pipe.api.event.Event;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
+import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -71,6 +75,7 @@ import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.util.Comparator;
 import java.util.Iterator;
+import java.util.List;
 import java.util.Queue;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.LinkedBlockingDeque;
@@ -109,6 +114,9 @@ public class IoTConsensusV2AsyncSink extends IoTDBSink 
implements ConsensusPipeS
   private ScheduledExecutorService backgroundTaskService;
   private final AtomicBoolean isClosed = new AtomicBoolean(false);
   private final int thisDataNodeId = 
IoTDBDescriptor.getInstance().getConfig().getDataNodeId();
+  private final TEndPoint localEndPoint =
+      new TEndPoint(IOTDB_CONFIG.getInternalAddress(), 
IOTDB_CONFIG.getDataRegionConsensusPort());
+  private DataRegionId dataRegionId;
   private IoTConsensusV2SinkMetrics iotConsensusV2SinkMetrics;
   private String consensusPipeName;
   private int consensusGroupId;
@@ -138,6 +146,7 @@ public class IoTConsensusV2AsyncSink extends IoTDBSink 
implements ConsensusPipeS
 
     // Get consensusGroupId from parameters passed by IoTConsensusV2Impl
     consensusGroupId = parameters.getInt(CONNECTOR_CONSENSUS_GROUP_ID_KEY);
+    dataRegionId = new DataRegionId(consensusGroupId);
     // Get consensusPipeName from parameters passed by IoTConsensusV2Impl
     consensusPipeName = parameters.getString(CONNECTOR_CONSENSUS_PIPE_NAME);
 
@@ -740,12 +749,47 @@ public class IoTConsensusV2AsyncSink extends IoTDBSink 
implements ConsensusPipeS
     }
   }
 
-  private TEndPoint getFollowerUrl() {
+  public TEndPoint getFollowerUrl() {
     // In current iotConsensusV2 design, one connector corresponds to one 
follower, so the peers is
     // actually a singleton list
     return nodeUrls.get(0);
   }
 
+  public boolean isUserDataTransferAuditEnabled() {
+    return dataRegionId != null && 
DataNodeUserDataTransferAuditor.isEnabledFor(dataRegionId);
+  }
+
+  public void recordUserDataTransferAudit(boolean success, String errorCode, 
Throwable error) {
+    if (!isUserDataTransferAuditEnabled()) {
+      return;
+    }
+    recordUserDataTransferAuditWithoutGroupCheck(success, errorCode, error);
+  }
+
+  public boolean recordUserDataTransferAudit(List<TSStatus> statuses) {
+    if (!isUserDataTransferAuditEnabled()) {
+      return false;
+    }
+    // The batch RPC is one physical transfer attempt. Keep one representative 
error value in the
+    // minimum audit record instead of concatenating an unbounded number of 
response details.
+    final TSStatus firstFailedStatus =
+        statuses.stream()
+            .filter(status -> status.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode())
+            .findFirst()
+            .orElse(null);
+    recordUserDataTransferAuditWithoutGroupCheck(
+        firstFailedStatus == null,
+        firstFailedStatus == null ? null : 
String.valueOf(firstFailedStatus.getCode()),
+        null);
+    return true;
+  }
+
+  private void recordUserDataTransferAuditWithoutGroupCheck(
+      boolean success, String errorCode, Throwable error) {
+    DataNodeUserDataTransferAuditor.record(
+        localEndPoint, localEndPoint, getFollowerUrl(), success, errorCode, 
error);
+  }
+
   // synchronized to avoid close connector when transfer event
   @Override
   public synchronized void close() {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2SyncSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2SyncSink.java
index 2f274189809..02b8f7cd15d 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2SyncSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/IoTConsensusV2SyncSink.java
@@ -25,6 +25,7 @@ import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.client.sync.SyncIoTConsensusV2ServiceClient;
+import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import 
org.apache.iotdb.commons.consensus.iotv2.container.IoTV2GlobalComponentContainer;
 import 
org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkRetryTimesConfigurableException;
@@ -36,6 +37,8 @@ import 
org.apache.iotdb.consensus.iotconsensusv2.thrift.TIoTConsensusV2BatchTran
 import 
org.apache.iotdb.consensus.iotconsensusv2.thrift.TIoTConsensusV2BatchTransferResp;
 import 
org.apache.iotdb.consensus.iotconsensusv2.thrift.TIoTConsensusV2TransferReq;
 import 
org.apache.iotdb.consensus.iotconsensusv2.thrift.TIoTConsensusV2TransferResp;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodePipeMessages;
 import org.apache.iotdb.db.pipe.consensus.metric.IoTConsensusV2SinkMetrics;
 import org.apache.iotdb.db.pipe.event.common.deletion.PipeDeleteDataNodeEvent;
@@ -87,6 +90,8 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
   private final int thisDataNodeId;
   private final int consensusGroupId;
   private final IoTConsensusV2SinkMetrics iotConsensusV2SinkMetrics;
+  private final TEndPoint localEndPoint;
+  private final DataRegionId dataRegionId;
   private IoTConsensusV2SyncBatchReqBuilder tabletBatchBuilder;
 
   public IoTConsensusV2SyncSink(
@@ -99,7 +104,12 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
     // retain the implementation of list to cope with possible future expansion
     this.peers = peers;
     this.consensusGroupId = consensusGroupId;
+    this.dataRegionId = new DataRegionId(consensusGroupId);
     this.thisDataNodeId = thisDataNodeId;
+    this.localEndPoint =
+        new TEndPoint(
+            IoTDBDescriptor.getInstance().getConfig().getInternalAddress(),
+            
IoTDBDescriptor.getInstance().getConfig().getDataRegionConsensusPort());
     this.syncRetryClientManager =
         
IoTV2GlobalComponentContainer.getInstance().getGlobalSyncClientManager();
     this.iotConsensusV2SinkMetrics = iotConsensusV2SinkMetrics;
@@ -199,6 +209,7 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
   }
 
   private void doTransfer() {
+    boolean transferAttemptRecorded = false;
     try (final SyncIoTConsensusV2ServiceClient syncIoTConsensusV2ServiceClient 
=
         syncRetryClientManager.borrowClient(getFollowerUrl())) {
       final TIoTConsensusV2BatchTransferResp resp;
@@ -210,6 +221,15 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
           resp.getBatchResps().stream()
               .map(TIoTConsensusV2TransferResp::getStatus)
               .collect(Collectors.toList());
+      if (isUserDataTransferAuditEnabled()) {
+        final TSStatus failedStatus =
+            statusList.stream().filter(status -> 
!isSuccessful(status)).findFirst().orElse(null);
+        recordTransferAttemptWithoutGroupCheck(
+            failedStatus == null,
+            failedStatus == null ? null : 
String.valueOf(failedStatus.getCode()),
+            null);
+        transferAttemptRecorded = true;
+      }
 
       // TODO(support batch): handle retry logic
       // Only handle the failed statuses to avoid string format performance 
overhead
@@ -225,6 +245,9 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
 
       tabletBatchBuilder.onSuccess();
     } catch (final Exception e) {
+      if (!transferAttemptRecorded) {
+        recordTransferAttempt(false, null, e);
+      }
       throw new PipeRuntimeSinkRetryTimesConfigurableException(
           String.format(
               IOT_CONSENSUS_V2_SYNC_CONNECTION_FAILED_FORMAT,
@@ -332,6 +355,7 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
             pipeInsertNodeTabletInsertionEvent.getRebootTimes());
     final TConsensusGroupId tConsensusGroupId =
         new TConsensusGroupId(TConsensusGroupType.DataRegion, 
consensusGroupId);
+    boolean transferAttemptRecorded = false;
 
     try (final SyncIoTConsensusV2ServiceClient syncIoTConsensusV2ServiceClient 
=
         syncRetryClientManager.borrowClient(getFollowerUrl())) {
@@ -342,7 +366,16 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
           IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq(
               insertNode, tCommitId, tConsensusGroupId, progressIndex, 
thisDataNodeId);
       resp = syncIoTConsensusV2ServiceClient.iotConsensusV2Transfer(req);
+      final TSStatus status = resp.getStatus();
+      recordTransferAttempt(
+          isSuccessful(status),
+          isSuccessful(status) ? null : String.valueOf(status.getCode()),
+          null);
+      transferAttemptRecorded = true;
     } catch (final Exception e) {
+      if (!transferAttemptRecorded) {
+        recordTransferAttempt(false, null, e);
+      }
       throw new PipeRuntimeSinkRetryTimesConfigurableException(
           String.format(
               IOT_CONSENSUS_V2_SYNC_CONNECTION_FAILED_FORMAT,
@@ -467,6 +500,8 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
                 ? readBuffer
                 : Arrays.copyOfRange(readBuffer, 0, readLength);
         final IoTConsensusV2TransferFilePieceResp resp;
+        final long transferPosition = position;
+        boolean transferAttemptRecorded = false;
         try {
           resp =
               
IoTConsensusV2TransferFilePieceResp.fromTIoTConsensusV2TransferResp(
@@ -486,7 +521,16 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
                               tCommitId,
                               tConsensusGroupId,
                               thisDataNodeId)));
+          final TSStatus transferStatus = resp.getStatus();
+          recordTransferAttempt(
+              isSuccessful(transferStatus),
+              isSuccessful(transferStatus) ? null : 
String.valueOf(transferStatus.getCode()),
+              null);
+          transferAttemptRecorded = true;
         } catch (Exception e) {
+          if (!transferAttemptRecorded) {
+            recordTransferAttempt(false, null, e);
+          }
           throw new PipeRuntimeSinkRetryTimesConfigurableException(
               String.format(
                   DataNodePipeMessages
@@ -535,6 +579,28 @@ public class IoTConsensusV2SyncSink extends IoTDBSink {
     return peers.get(0);
   }
 
+  private static boolean isSuccessful(TSStatus status) {
+    return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+        || status.getCode() == 
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode();
+  }
+
+  private void recordTransferAttempt(boolean success, String errorCode, 
Throwable error) {
+    if (!isUserDataTransferAuditEnabled()) {
+      return;
+    }
+    recordTransferAttemptWithoutGroupCheck(success, errorCode, error);
+  }
+
+  private void recordTransferAttemptWithoutGroupCheck(
+      boolean success, String errorCode, Throwable error) {
+    DataNodeUserDataTransferAuditor.record(
+        localEndPoint, localEndPoint, getFollowerUrl(), success, errorCode, 
error);
+  }
+
+  private boolean isUserDataTransferAuditEnabled() {
+    return DataNodeUserDataTransferAuditor.isEnabledFor(dataRegionId);
+  }
+
   // synchronized to avoid close connector when transfer event
   @Override
   public synchronized void close() {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletBatchEventHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletBatchEventHandler.java
index 54d8b3e6319..527fca89001 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletBatchEventHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletBatchEventHandler.java
@@ -52,6 +52,7 @@ public class IoTConsensusV2TabletBatchEventHandler
   private final TIoTConsensusV2BatchTransferReq req;
   private final IoTConsensusV2AsyncSink connector;
   private final IoTConsensusV2SinkMetrics iotConsensusV2SinkMetrics;
+  private boolean transferAuditRecorded;
 
   public IoTConsensusV2TabletBatchEventHandler(
       final IoTConsensusV2AsyncBatchReqBuilder batchBuilder,
@@ -84,6 +85,7 @@ public class IoTConsensusV2TabletBatchEventHandler
           response.getBatchResps().stream()
               .map(TIoTConsensusV2TransferResp::getStatus)
               .collect(Collectors.toList());
+      transferAuditRecorded = connector.recordUserDataTransferAudit(status);
 
       if (status.stream()
           .anyMatch(
@@ -118,6 +120,12 @@ public class IoTConsensusV2TabletBatchEventHandler
 
   @Override
   public void onError(final Exception exception) {
+    // A retry is sent through a new handler and produces its own audit event. 
This guard only
+    // prevents a post-response processing exception from recording the same 
attempt twice.
+    if (!transferAuditRecorded) {
+      connector.recordUserDataTransferAudit(false, null, exception);
+      transferAuditRecorded = true;
+    }
     final Object pipeNames =
         events.stream()
             .map(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletInsertionEventHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletInsertionEventHandler.java
index 3e1efd17088..794414c9d30 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletInsertionEventHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TabletInsertionEventHandler.java
@@ -54,6 +54,7 @@ public abstract class 
IoTConsensusV2TabletInsertionEventHandler<
   protected final IoTConsensusV2SinkMetrics metric;
 
   private final long createTime;
+  private boolean transferAuditRecorded;
 
   protected IoTConsensusV2TabletInsertionEventHandler(
       TabletInsertionEvent event,
@@ -83,6 +84,12 @@ public abstract class 
IoTConsensusV2TabletInsertionEventHandler<
     }
 
     final TSStatus status = response.getStatus();
+    final boolean success =
+        status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+            || status.getCode() == 
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode();
+    connector.recordUserDataTransferAudit(
+        success, success ? null : String.valueOf(status.getCode()), null);
+    transferAuditRecorded = true;
     try {
       // Only handle the failed statuses to avoid string format performance 
overhead
       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()
@@ -113,6 +120,10 @@ public abstract class 
IoTConsensusV2TabletInsertionEventHandler<
 
   @Override
   public void onError(Exception exception) {
+    if (!transferAuditRecorded) {
+      connector.recordUserDataTransferAudit(false, null, exception);
+      transferAuditRecorded = true;
+    }
     EnrichedEvent event = (EnrichedEvent) this.event;
     PipeLogger.log(
         ignored ->
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TsFileInsertionEventHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TsFileInsertionEventHandler.java
index 8ebff374531..f350ff8c78f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TsFileInsertionEventHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/handler/IoTConsensusV2TsFileInsertionEventHandler.java
@@ -87,6 +87,8 @@ public class IoTConsensusV2TsFileInsertionEventHandler
   private final long createTime;
 
   private long startTransferPieceTime;
+  private boolean currentAttemptContainsUserData;
+  private boolean transferAuditRecorded;
 
   public IoTConsensusV2TsFileInsertionEventHandler(
       final PipeTsFileInsertionEvent event,
@@ -161,6 +163,7 @@ public class IoTConsensusV2TsFileInsertionEventHandler
         transfer(client);
       } else if (currentFile == tsFile) {
         isSealSignalSent.set(true);
+        currentAttemptContainsUserData = false;
         client.iotConsensusV2Transfer(
             transferMod
                 ? 
IoTConsensusV2TsFileSealWithModReq.toTIoTConsensusV2TransferReq(
@@ -191,6 +194,8 @@ public class IoTConsensusV2TsFileInsertionEventHandler
         readLength == readFileBufferSize
             ? readBuffer
             : Arrays.copyOfRange(readBuffer, 0, readLength);
+    currentAttemptContainsUserData = true;
+    transferAuditRecorded = false;
     client.iotConsensusV2Transfer(
         transferMod
             ? IoTConsensusV2TsFilePieceWithModReq.toTIoTConsensusV2TransferReq(
@@ -275,6 +280,13 @@ public class IoTConsensusV2TsFileInsertionEventHandler
     try {
       final IoTConsensusV2TransferFilePieceResp resp =
           
IoTConsensusV2TransferFilePieceResp.fromTIoTConsensusV2TransferResp(response);
+      final TSStatus transferStatus = resp.getStatus();
+      final boolean success =
+          transferStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()
+              || transferStatus.getCode() == 
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode();
+      connector.recordUserDataTransferAudit(
+          success, success ? null : String.valueOf(transferStatus.getCode()), 
null);
+      transferAuditRecorded = true;
 
       // This case only happens when the connection is broken, and the 
connector is reconnected
       // to the receiver, then the receiver will redirect the file position to 
the last position
@@ -309,6 +321,10 @@ public class IoTConsensusV2TsFileInsertionEventHandler
 
   @Override
   public void onError(final Exception exception) {
+    if (currentAttemptContainsUserData && !transferAuditRecorded) {
+      connector.recordUserDataTransferAudit(false, null, exception);
+      transferAuditRecorded = true;
+    }
     PipeLogger.log(
         ignored ->
             LOGGER.warn(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
index 2cd33b2d769..67a5defb09a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
@@ -20,9 +20,11 @@
 package org.apache.iotdb.db.queryengine.execution.exchange.source;
 
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.commons.audit.UserDataTransferErrorCode;
 import org.apache.iotdb.commons.client.IClientManager;
 import 
org.apache.iotdb.commons.client.sync.SyncDataNodeMPPDataExchangeServiceClient;
 import org.apache.iotdb.commons.utils.TestOnly;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
 import org.apache.iotdb.db.queryengine.common.FragmentInstanceId;
@@ -40,6 +42,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TGetDataBlockResponse;
 
 import com.google.common.util.concurrent.ListenableFuture;
 import com.google.common.util.concurrent.SettableFuture;
+import org.apache.thrift.TException;
 import org.apache.tsfile.external.commons.lang3.Validate;
 import org.apache.tsfile.read.common.block.TsBlock;
 import org.apache.tsfile.read.common.block.column.TsBlockSerde;
@@ -73,6 +76,7 @@ public class SourceHandle implements ISourceHandle {
   private static final long DEFAULT_RETRY_INTERVAL_IN_MS = 1000;
 
   private final TEndPoint remoteEndpoint;
+  private final TEndPoint localEndpoint;
   private final TFragmentInstanceId remoteFragmentInstanceId;
   private final TFragmentInstanceId localFragmentInstanceId;
 
@@ -175,6 +179,10 @@ public class SourceHandle implements ISourceHandle {
         Validate.notNull(
             remoteEndpoint,
             
DataNodeQueryMessages.EXCEPTION_REMOTEENDPOINT_CAN_NOT_BE_NULL_DOT_DE2B5885);
+    this.localEndpoint =
+        new TEndPoint(
+            IoTDBDescriptor.getInstance().getConfig().getInternalAddress(),
+            
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangePort());
     this.remoteFragmentInstanceId =
         Validate.notNull(
             remoteFragmentInstanceId,
@@ -637,11 +645,19 @@ public class SourceHandle implements ISourceHandle {
           attempt += 1;
 
           long startTime = System.nanoTime();
+          boolean transferAttemptRecorded = false;
           try (SyncDataNodeMPPDataExchangeServiceClient client =
               
mppDataExchangeServiceClientManager.borrowClient(remoteEndpoint)) {
             TGetDataBlockResponse resp = client.getDataBlock(req);
             int tsBlockNum = resp.getTsBlocks().size();
-            if (tsBlockNum == 0) {
+            if (tsBlockNum != endSequenceId - startSequenceId) {
+              recordTransferAttempt(
+                  false,
+                  tsBlockNum == 0
+                      ? UserDataTransferErrorCode.EMPTY_RESPONSE.name()
+                      : 
UserDataTransferErrorCode.UNEXPECTED_RESPONSE_SIZE.name(),
+                  null);
+              transferAttemptRecorded = true;
               if (!closed) {
                 // failed to pull TsBlocks
                 LOGGER.warn(
@@ -652,7 +668,11 @@ public class SourceHandle implements ISourceHandle {
                     remoteFragmentInstanceId,
                     indexOfUpstreamSinkHandle);
               }
-              return;
+              if (tsBlockNum == 0) {
+                return;
+              }
+              throw new TException(
+                  
DataNodeQueryMessages.EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33);
             }
             List<ByteBuffer> tsBlocks = new ArrayList<>(tsBlockNum);
             tsBlocks.addAll(resp.getTsBlocks());
@@ -664,23 +684,37 @@ public class SourceHandle implements ISourceHandle {
                 GET_DATA_BLOCK_NUM_CALLER, tsBlockNum);
             executorService.submit(
                 new SendAcknowledgeDataBlockEventTask(startSequenceId, 
endSequenceId));
+            boolean receiverClosed = false;
             synchronized (SourceHandle.this) {
               if (aborted || closed) {
-                return;
-              }
-              for (int i = startSequenceId; i < endSequenceId; i++) {
-                sequenceIdToTsBlock.put(i, tsBlocks.get(i - startSequenceId));
-              }
-              if (LOGGER.isDebugEnabled()) {
-                LOGGER.debug(DataNodeQueryMessages.PUT_TSBLOCKS_INTO_BUFFER);
-              }
-              if (!blocked.isDone()) {
-                blocked.set(null);
+                receiverClosed = true;
+              } else {
+                for (int i = startSequenceId; i < endSequenceId; i++) {
+                  sequenceIdToTsBlock.put(i, tsBlocks.get(i - 
startSequenceId));
+                }
+                if (LOGGER.isDebugEnabled()) {
+                  LOGGER.debug(DataNodeQueryMessages.PUT_TSBLOCKS_INTO_BUFFER);
+                }
+                if (!blocked.isDone()) {
+                  blocked.set(null);
+                }
               }
             }
+            recordTransferAttempt(
+                !receiverClosed,
+                receiverClosed ? 
UserDataTransferErrorCode.RECEIVER_CLOSED.name() : null,
+                null);
+            transferAttemptRecorded = true;
+            if (receiverClosed) {
+              return;
+            }
             break;
           } catch (Throwable e) {
 
+            if (!transferAttemptRecorded) {
+              recordTransferAttempt(false, null, e);
+            }
+
             LOGGER.warn(
                 DataNodeQueryMessages.FAILED_TO_GET_DATA_BLOCK,
                 startSequenceId,
@@ -710,6 +744,11 @@ public class SourceHandle implements ISourceHandle {
       }
     }
 
+    private void recordTransferAttempt(boolean success, String errorCode, 
Throwable error) {
+      DataNodeUserDataTransferAuditor.record(
+          localEndpoint, remoteEndpoint, localEndpoint, success, errorCode, 
error);
+    }
+
     private void fail(Throwable t) {
       synchronized (SourceHandle.this) {
         if (aborted || closed) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSender.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSender.java
index da5a2dfc889..60a1739720e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSender.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSender.java
@@ -21,8 +21,12 @@ package org.apache.iotdb.db.queryengine.plan.scheduler;
 
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.auth.entity.User;
 import org.apache.iotdb.commons.client.IClientManager;
 import 
org.apache.iotdb.commons.client.async.AsyncDataNodeInternalServiceClient;
+import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.FragmentInstance;
 import org.apache.iotdb.mpp.rpc.thrift.TPlanNode;
@@ -49,6 +53,7 @@ public class AsyncPlanNodeSender {
   private final IClientManager<TEndPoint, AsyncDataNodeInternalServiceClient>
       asyncInternalServiceClientManager;
   private final List<FragmentInstance> instances;
+  private final TEndPoint localEndPoint;
 
   private final Map<TEndPoint, BatchRequestWithIndex> batchRequests;
   private final Map<Integer, TSendSinglePlanNodeResp> instanceId2RespMap;
@@ -65,6 +70,10 @@ public class AsyncPlanNodeSender {
     this.startSendTime = System.nanoTime();
     this.asyncInternalServiceClientManager = asyncInternalServiceClientManager;
     this.instances = instances;
+    this.localEndPoint =
+        new TEndPoint(
+            IoTDBDescriptor.getInstance().getConfig().getInternalAddress(),
+            IoTDBDescriptor.getInstance().getConfig().getInternalPort());
     this.batchRequests = new HashMap<>();
     for (int i = 0; i < instances.size(); i++) {
       this.batchRequests
@@ -76,7 +85,8 @@ public class AsyncPlanNodeSender {
               new TSendSinglePlanNodeReq(
                   new TPlanNode(
                       
instances.get(i).getFragment().getPlanNodeTree().serializeToByteBuffer()),
-                  instances.get(i).getRegionReplicaSet().getRegionId()));
+                  instances.get(i).getRegionReplicaSet().getRegionId()),
+              containsUserData(instances.get(i)));
     }
     this.instanceId2RespMap = new ConcurrentHashMap<>(instances.size() + 1, 1);
     this.needRetryInstanceIndex = Collections.synchronizedList(new 
ArrayList<>());
@@ -91,7 +101,10 @@ public class AsyncPlanNodeSender {
               pendingNumber,
               instanceId2RespMap,
               needRetryInstanceIndex,
-              startSendTime);
+              startSendTime,
+              localEndPoint,
+              entry.getKey(),
+              entry.getValue().containsUserData());
       try {
         AsyncDataNodeInternalServiceClient client =
             asyncInternalServiceClientManager.borrowClient(entry.getKey());
@@ -188,7 +201,8 @@ public class AsyncPlanNodeSender {
                           .getFragment()
                           .getPlanNodeTree()
                           .serializeToByteBuffer()),
-                  
instances.get(fragmentInstanceIndex).getRegionReplicaSet().getRegionId()));
+                  
instances.get(fragmentInstanceIndex).getRegionReplicaSet().getRegionId()),
+              containsUserData(instances.get(fragmentInstanceIndex)));
     }
 
     // 2. reset the pendingNumber, needRetryInstanceIds and startSendTime
@@ -213,9 +227,13 @@ public class AsyncPlanNodeSender {
     private final List<Integer> indexes = new ArrayList<>();
     private final TSendBatchPlanNodeReq batchRequest = new 
TSendBatchPlanNodeReq();
 
-    void addSinglePlanNodeReq(int index, TSendSinglePlanNodeReq singleRequest) 
{
+    private boolean containsUserData;
+
+    void addSinglePlanNodeReq(
+        int index, TSendSinglePlanNodeReq singleRequest, boolean 
containsUserData) {
       indexes.add(index);
       batchRequest.addToRequests(singleRequest);
+      this.containsUserData |= containsUserData;
     }
 
     public List<Integer> getIndexes() {
@@ -225,5 +243,25 @@ public class AsyncPlanNodeSender {
     public TSendBatchPlanNodeReq getBatchRequest() {
       return batchRequest;
     }
+
+    public boolean containsUserData() {
+      return containsUserData;
+    }
+  }
+
+  private static boolean containsUserData(FragmentInstance instance) {
+    return containsUserData(
+        instance.getFragment().getPlanNodeTree(),
+        instance.getSessionInfo() == null ? null : 
instance.getSessionInfo().getUserName());
+  }
+
+  static boolean containsUserData(PlanNode node, String username) {
+    return DataNodeUserDataTransferAuditor.isEnabled()
+        && !User.BUILTIN_INTERNAL_AUDIT_LOG_USERNAME.equals(username)
+        && DataNodeUserDataTransferAuditor.containsUserData(node);
+  }
+
+  static boolean containsUserData(PlanNode node) {
+    return DataNodeUserDataTransferAuditor.containsUserData(node);
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncSendPlanNodeHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncSendPlanNodeHandler.java
index 23f99da3484..eae8dbc0091 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncSendPlanNodeHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncSendPlanNodeHandler.java
@@ -19,8 +19,11 @@
 
 package org.apache.iotdb.db.queryengine.plan.scheduler;
 
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.commons.audit.UserDataTransferErrorCode;
 import org.apache.iotdb.commons.service.metric.PerformanceOverviewMetrics;
 import org.apache.iotdb.commons.utils.StatusUtils;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.mpp.rpc.thrift.TSendBatchPlanNodeResp;
 import org.apache.iotdb.mpp.rpc.thrift.TSendSinglePlanNodeResp;
 import org.apache.iotdb.rpc.RpcUtils;
@@ -43,6 +46,10 @@ public class AsyncSendPlanNodeHandler implements 
AsyncMethodCallback<TSendBatchP
   private final Map<Integer, TSendSinglePlanNodeResp> instanceId2RespMap;
   private final List<Integer> needRetryInstanceIndex;
   private final long sendTime;
+  private final TEndPoint localEndPoint;
+  private final TEndPoint targetEndPoint;
+  private final boolean containsUserData;
+  private boolean transferAuditRecorded;
   private static final PerformanceOverviewMetrics PERFORMANCE_OVERVIEW_METRICS 
=
       PerformanceOverviewMetrics.getInstance();
 
@@ -51,16 +58,23 @@ public class AsyncSendPlanNodeHandler implements 
AsyncMethodCallback<TSendBatchP
       AtomicLong pendingNumber,
       Map<Integer, TSendSinglePlanNodeResp> instanceId2RespMap,
       List<Integer> needRetryInstanceIndex,
-      long sendTime) {
+      long sendTime,
+      TEndPoint localEndPoint,
+      TEndPoint targetEndPoint,
+      boolean containsUserData) {
     this.instanceIds = instanceIds;
     this.pendingNumber = pendingNumber;
     this.instanceId2RespMap = instanceId2RespMap;
     this.needRetryInstanceIndex = needRetryInstanceIndex;
     this.sendTime = sendTime;
+    this.localEndPoint = localEndPoint;
+    this.targetEndPoint = targetEndPoint;
+    this.containsUserData = containsUserData;
   }
 
   @Override
   public void onComplete(TSendBatchPlanNodeResp sendBatchPlanNodeResp) {
+    recordTransferAttempt(sendBatchPlanNodeResp);
     for (int i = 0; i < sendBatchPlanNodeResp.getResponses().size(); i++) {
       TSendSinglePlanNodeResp singlePlanNodeResp = 
sendBatchPlanNodeResp.getResponses().get(i);
       instanceId2RespMap.put(instanceIds.get(i), singlePlanNodeResp);
@@ -78,6 +92,9 @@ public class AsyncSendPlanNodeHandler implements 
AsyncMethodCallback<TSendBatchP
 
   @Override
   public void onError(Exception e) {
+    if (!transferAuditRecorded) {
+      recordTransferAttempt(false, null, e);
+    }
     if (needRetry(e)) {
       needRetryInstanceIndex.addAll(instanceIds);
     }
@@ -106,4 +123,38 @@ public class AsyncSendPlanNodeHandler implements 
AsyncMethodCallback<TSendBatchP
   private boolean needRetry(TSendSinglePlanNodeResp resp) {
     return !resp.accepted && resp.status != null && 
StatusUtils.needRetryHelper(resp.status);
   }
+
+  private void recordTransferAttempt(TSendBatchPlanNodeResp response) {
+    if (!containsUserData) {
+      return;
+    }
+    if (response.getResponsesSize() != instanceIds.size()) {
+      recordTransferAttempt(false, 
UserDataTransferErrorCode.REMOTE_REJECTED.name(), null);
+      return;
+    }
+    for (TSendSinglePlanNodeResp singleResponse : response.getResponses()) {
+      if (!singleResponse.isAccepted()
+          || (singleResponse.isSetStatus()
+              && singleResponse.getStatus().getCode()
+                  != TSStatusCode.SUCCESS_STATUS.getStatusCode())) {
+        recordTransferAttempt(
+            false,
+            singleResponse.isSetStatus()
+                ? String.valueOf(singleResponse.getStatus().getCode())
+                : UserDataTransferErrorCode.REMOTE_REJECTED.name(),
+            null);
+        return;
+      }
+    }
+    recordTransferAttempt(true, null, null);
+  }
+
+  private void recordTransferAttempt(boolean success, String errorCode, 
Throwable error) {
+    if (!containsUserData) {
+      return;
+    }
+    DataNodeUserDataTransferAuditor.record(
+        localEndPoint, localEndPoint, targetEndPoint, success, errorCode, 
error);
+    transferAuditRecorded = true;
+  }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
index 7ae9d8d6f77..24c42a0c16f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.calc.metric.QueryExecutionMetricSet;
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.audit.UserDataTransferErrorCode;
 import org.apache.iotdb.commons.client.IClientManager;
 import 
org.apache.iotdb.commons.client.async.AsyncDataNodeInternalServiceClient;
 import org.apache.iotdb.commons.client.exception.ClientManagerException;
@@ -35,6 +36,7 @@ import 
org.apache.iotdb.commons.service.metric.PerformanceOverviewMetrics;
 import org.apache.iotdb.commons.utils.RetryUtils;
 import org.apache.iotdb.consensus.exception.ConsensusGroupNotExistException;
 import org.apache.iotdb.consensus.exception.RatisReadUnavailableException;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.mpp.FragmentInstanceDispatchException;
 import org.apache.iotdb.db.exception.query.QueryTimeoutRuntimeException;
@@ -471,6 +473,12 @@ public class FragmentInstanceDispatcherImpl implements 
IFragInstanceDispatcher {
           ClientManagerException,
           RatisReadUnavailableException,
           ConsensusGroupNotExistException {
+    final boolean containsUserData =
+        (instance.getType() == QueryType.WRITE || instance.getType() == 
QueryType.OTHER)
+            && AsyncPlanNodeSender.containsUserData(
+                instance.getFragment().getPlanNodeTree(),
+                instance.getSessionInfo() == null ? null : 
instance.getSessionInfo().getUserName());
+    boolean transferAttemptRecorded = false;
     try (final SyncDataNodeInternalServiceClient client =
         syncInternalServiceClientManager.borrowClient(endPoint)) {
       switch (instance.getType()) {
@@ -523,6 +531,25 @@ public class FragmentInstanceDispatcherImpl implements 
IFragInstanceDispatcher {
                           instance.getRegionReplicaSet().getRegionId())));
           final TSendSinglePlanNodeResp sendPlanNodeResp =
               client.sendBatchPlanNode(sendPlanNodeReq).getResponses().get(0);
+          if (containsUserData) {
+            final boolean success =
+                sendPlanNodeResp.isAccepted()
+                    && (!sendPlanNodeResp.isSetStatus()
+                        || sendPlanNodeResp.getStatus().getCode()
+                            == TSStatusCode.SUCCESS_STATUS.getStatusCode());
+            DataNodeUserDataTransferAuditor.record(
+                new TEndPoint(localhostIpAddr, localhostInternalPort),
+                new TEndPoint(localhostIpAddr, localhostInternalPort),
+                endPoint,
+                success,
+                success
+                    ? null
+                    : sendPlanNodeResp.isSetStatus()
+                        ? 
String.valueOf(sendPlanNodeResp.getStatus().getCode())
+                        : UserDataTransferErrorCode.REMOTE_REJECTED.name(),
+                null);
+            transferAttemptRecorded = true;
+          }
           if (!sendPlanNodeResp.accepted) {
             if (sendPlanNodeResp.getStatus() == null) {
               throw new FragmentInstanceDispatchException(
@@ -556,7 +583,15 @@ public class FragmentInstanceDispatcherImpl implements 
IFragInstanceDispatcher {
                   TSStatusCode.EXECUTE_STATEMENT_ERROR,
                   String.format(DataNodeQueryMessages.UNKNOWN_READ_TYPE_FMT, 
instance.getType())));
       }
+    } catch (ClientManagerException e) {
+      if (!transferAttemptRecorded) {
+        recordWriteTransferFailureIfNecessary(endPoint, containsUserData, e);
+      }
+      throw e;
     } catch (TException e) {
+      if (!transferAttemptRecorded) {
+        recordWriteTransferFailureIfNecessary(endPoint, containsUserData, e);
+      }
       Throwable rootCause = ExceptionUtils.getRootCause(e);
       if (rootCause instanceof TTransportException
           && ((TTransportException) rootCause).getType() == 
TTransportException.CORRUPTED_DATA) {
@@ -568,6 +603,16 @@ public class FragmentInstanceDispatcherImpl implements 
IFragInstanceDispatcher {
     }
   }
 
+  private void recordWriteTransferFailureIfNecessary(
+      TEndPoint endPoint, boolean containsUserData, Throwable error) {
+    if (!containsUserData) {
+      return;
+    }
+    final TEndPoint localEndPoint = new TEndPoint(localhostIpAddr, 
localhostInternalPort);
+    DataNodeUserDataTransferAuditor.record(
+        localEndPoint, localEndPoint, endPoint, false, null, error);
+  }
+
   private void dispatchRemoteFailed(TEndPoint endPoint, Exception e)
       throws FragmentInstanceDispatchException {
     LOGGER.warn(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileDispatcherImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileDispatcherImpl.java
index 3f5020cccf3..1ba39ffd0bd 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileDispatcherImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileDispatcherImpl.java
@@ -24,6 +24,7 @@ import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
+import org.apache.iotdb.commons.audit.UserDataTransferErrorCode;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.client.sync.SyncDataNodeInternalServiceClient;
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
@@ -32,6 +33,7 @@ import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import org.apache.iotdb.commons.consensus.index.ProgressIndexType;
 import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.audit.DataNodeUserDataTransferAuditor;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.load.LoadFileException;
 import org.apache.iotdb.db.exception.mpp.FragmentInstanceDispatchException;
@@ -222,16 +224,30 @@ public class LoadTsFileDispatcherImpl implements 
IFragInstanceDispatcher, AutoCl
 
   private void dispatchRemote(TTsFilePieceReq loadTsFileReq, TEndPoint 
endPoint)
       throws FragmentInstanceDispatchException {
+    boolean transferAttemptRecorded = false;
     try (SyncDataNodeInternalServiceClient client =
         internalServiceClientManager.borrowClient(endPoint)) {
       client.setTimeout(CONNECTION_TIMEOUT_MS.get());
 
       final TLoadResp loadResp = client.sendTsFilePieceNode(loadTsFileReq);
       if (!loadResp.isAccepted()) {
+        recordTransferAttempt(
+            endPoint,
+            false,
+            loadResp.isSetStatus()
+                ? String.valueOf(loadResp.getStatus().getCode())
+                : UserDataTransferErrorCode.REMOTE_REJECTED.name(),
+            null);
+        transferAttemptRecorded = true;
         LOGGER.warn(loadResp.message);
         throw new FragmentInstanceDispatchException(loadResp.status);
       }
+      recordTransferAttempt(endPoint, true, null, null);
+      transferAttemptRecorded = true;
     } catch (Exception e) {
+      if (!transferAttemptRecorded) {
+        recordTransferAttempt(endPoint, false, null, e);
+      }
       adjustTimeoutIfNecessary(e);
 
       final String exceptionMessage =
@@ -246,6 +262,13 @@ public class LoadTsFileDispatcherImpl implements 
IFragInstanceDispatcher, AutoCl
     }
   }
 
+  private void recordTransferAttempt(
+      TEndPoint target, boolean success, String errorCode, Throwable error) {
+    final TEndPoint localEndPoint = new TEndPoint(localhostIpAddr, 
localhostInternalPort);
+    DataNodeUserDataTransferAuditor.record(
+        localEndPoint, localEndPoint, target, success, errorCode, error);
+  }
+
   public Future<FragInstanceDispatchResult> dispatchCommand(
       TLoadCommandReq originalLoadCommandReq, Set<TRegionReplicaSet> 
replicaSets) {
     Set<TEndPoint> allEndPoint = new HashSet<>();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
index 03438b60707..0e66509d350 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
@@ -766,6 +766,7 @@ public class WALNode implements IWALNode {
       AtomicLong currentEntryLocalSeq = new AtomicLong(-1);
       AtomicLong currentEntryPhysicalTime = new AtomicLong(0);
       AtomicLong currentEntryNodeId = new AtomicLong(-1);
+      AtomicBoolean currentEntryContainsUserData = new AtomicBoolean(false);
 
       long memorySize = 0;
 
@@ -779,10 +780,12 @@ public class WALNode implements IWALNode {
                   (localSeq >= 0)
                       ? new IndexedConsensusRequest(nextSearchIndex, localSeq, 
tmpNodes.get())
                       : new IndexedConsensusRequest(nextSearchIndex, 
tmpNodes.get());
-              req.setPhysicalTime(currentEntryPhysicalTime.get())
+              req.setContainsUserData(currentEntryContainsUserData.get())
+                  .setPhysicalTime(currentEntryPhysicalTime.get())
                   .setNodeId((int) currentEntryNodeId.get());
               insertNodes.add(req);
               tmpNodes.set(new ArrayList<>());
+              currentEntryContainsUserData.set(false);
               nextSearchIndex++;
               if (notFirstFile.get()) {
                 hasCollectedSufficientData.set(true);
@@ -819,6 +822,8 @@ public class WALNode implements IWALNode {
                 
currentEntryLocalSeq.set(walByteBufReader.getCurrentEntryLocalSeq());
                 
currentEntryPhysicalTime.set(walByteBufReader.getCurrentEntryPhysicalTime());
                 
currentEntryNodeId.set(walByteBufReader.getCurrentEntryNodeId());
+                currentEntryContainsUserData.set(
+                    currentEntryContainsUserData.get() || 
containsUserData(type));
                 if (type == WALEntryType.OBJECT_FILE_NODE) {
                   WALEntry walEntry =
                       WALEntry.deserialize(
@@ -854,6 +859,8 @@ public class WALNode implements IWALNode {
                 
currentEntryLocalSeq.set(walByteBufReader.getCurrentEntryLocalSeq());
                 
currentEntryPhysicalTime.set(walByteBufReader.getCurrentEntryPhysicalTime());
                 
currentEntryNodeId.set(walByteBufReader.getCurrentEntryNodeId());
+                currentEntryContainsUserData.set(
+                    currentEntryContainsUserData.get() || 
containsUserData(type));
                 if (type == WALEntryType.OBJECT_FILE_NODE) {
                   WALEntry walEntry =
                       WALEntry.deserialize(
@@ -909,6 +916,13 @@ public class WALNode implements IWALNode {
       return false;
     }
 
+    private boolean containsUserData(WALEntryType type) {
+      return type == WALEntryType.INSERT_ROW_NODE
+          || type == WALEntryType.INSERT_TABLET_NODE
+          || type == WALEntryType.INSERT_ROWS_NODE
+          || type == WALEntryType.OBJECT_FILE_NODE;
+    }
+
     @Override
     public IndexedConsensusRequest next() {
       if (itr == null && !hasNext()) {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditorTest.java
new file mode 100644
index 00000000000..0cba74868b6
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/audit/DataNodeUserDataTransferAuditorTest.java
@@ -0,0 +1,136 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.audit;
+
+import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.commons.request.IConsensusRequest;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.ObjectNode;
+import org.apache.iotdb.db.storageengine.StorageEngine;
+import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
+
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.powermock.api.mockito.PowerMockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
+import org.powermock.core.classloader.annotations.PrepareForTest;
+import org.powermock.modules.junit4.PowerMockRunner;
+
+import java.util.Collections;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@PowerMockIgnore({"com.sun.org.apache.xerces.*", "javax.xml.*", "org.xml.*", 
"javax.management.*"})
+@RunWith(PowerMockRunner.class)
+@PrepareForTest(StorageEngine.class)
+public class DataNodeUserDataTransferAuditorTest {
+
+  private final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
+  private boolean auditLogEnabled;
+
+  @Before
+  public void setUp() {
+    auditLogEnabled = commonConfig.isEnableAuditLog();
+  }
+
+  @After
+  public void tearDown() {
+    commonConfig.setEnableAuditLog(auditLogEnabled);
+  }
+
+  @Test
+  public void testAuditDatabaseIsExcludedFromGroupTransferAudit() {
+    assertFalse(DataNodeUserDataTransferAuditor.containsUserData("__audit"));
+    
assertFalse(DataNodeUserDataTransferAuditor.containsUserData("root.__audit"));
+    assertTrue(DataNodeUserDataTransferAuditor.containsUserData("root.sg"));
+  }
+
+  @Test
+  public void testIoTConsensusV2AuditGateExcludesAuditDatabase() {
+    final StorageEngine storageEngine = mock(StorageEngine.class);
+    final DataRegion auditDataRegion = mock(DataRegion.class);
+    final DataRegion userDataRegion = mock(DataRegion.class);
+    final DataRegionId auditDataRegionId = new DataRegionId(1);
+    final DataRegionId userDataRegionId = new DataRegionId(2);
+    PowerMockito.mockStatic(StorageEngine.class);
+    PowerMockito.when(StorageEngine.getInstance()).thenReturn(storageEngine);
+    
when(storageEngine.getDataRegion(auditDataRegionId)).thenReturn(auditDataRegion);
+    
when(storageEngine.getDataRegion(userDataRegionId)).thenReturn(userDataRegion);
+    when(auditDataRegion.getDatabaseName()).thenReturn("root.__audit");
+    when(userDataRegion.getDatabaseName()).thenReturn("root.sg");
+
+    commonConfig.setEnableAuditLog(true);
+
+    
assertFalse(DataNodeUserDataTransferAuditor.isEnabledFor(auditDataRegionId));
+    assertTrue(DataNodeUserDataTransferAuditor.isEnabledFor(userDataRegionId));
+  }
+
+  @Test
+  public void testIoTConsensusV2AuditGateShortCircuitsWhenDisabled() {
+    final StorageEngine storageEngine = mock(StorageEngine.class);
+    PowerMockito.mockStatic(StorageEngine.class);
+    PowerMockito.when(StorageEngine.getInstance()).thenReturn(storageEngine);
+    commonConfig.setEnableAuditLog(false);
+
+    assertFalse(DataNodeUserDataTransferAuditor.isEnabledFor(new 
DataRegionId(1)));
+
+    verify(storageEngine, never()).getDataRegion(new DataRegionId(1));
+  }
+
+  @Test
+  public void testAuditDatabaseIsExcludedFromConsensusTransferAudit() {
+    final InsertNode insertNode = mock(InsertNode.class);
+
+    assertFalse(DataNodeUserDataTransferAuditor.containsUserData("__audit", 
insertNode));
+    
assertFalse(DataNodeUserDataTransferAuditor.containsUserData("root.__audit", 
insertNode));
+    assertTrue(DataNodeUserDataTransferAuditor.containsUserData("root.sg", 
insertNode));
+  }
+
+  @Test
+  public void testNonInsertConsensusRequestIsExcluded() {
+    final PlanNode planNode = mock(PlanNode.class);
+    when(planNode.getChildren()).thenReturn(Collections.emptyList());
+
+    assertFalse(DataNodeUserDataTransferAuditor.containsUserData("root.sg", 
planNode));
+  }
+
+  @Test
+  public void testClassificationDoesNotDeserializeConsensusRequest() {
+    final IConsensusRequest request = mock(IConsensusRequest.class);
+
+    assertFalse(DataNodeUserDataTransferAuditor.containsUserData("root.sg", 
request));
+    verify(request, never()).serializeToByteBuffer();
+  }
+
+  @Test
+  public void testObjectFileNodeContainsUserData() {
+    
assertTrue(DataNodeUserDataTransferAuditor.containsUserData(mock(ObjectNode.class)));
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
index 66d50675ddf..03220a3d10b 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
@@ -610,6 +610,80 @@ public class SourceHandleTest {
         .onAborted(sourceHandle);
   }
 
+  @Test
+  public void testShortResponseRetriesAndFails() {
+    final String queryId = "q0";
+    final int numOfMockTsBlock = 10;
+    final TEndPoint remoteEndpoint =
+        new TEndPoint("remote", 
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangePort());
+    final TFragmentInstanceId remoteFragmentInstanceId = new 
TFragmentInstanceId(queryId, 1, "0");
+    final String localPlanNodeId = "exchange_0";
+    final TFragmentInstanceId localFragmentInstanceId = new 
TFragmentInstanceId(queryId, 0, "0");
+
+    final LocalMemoryManager mockLocalMemoryManager = 
Mockito.mock(LocalMemoryManager.class);
+    final MemoryPool mockMemoryPool = Utils.createMockNonBlockedMemoryPool();
+    
Mockito.when(mockLocalMemoryManager.getQueryPool()).thenReturn(mockMemoryPool);
+    final SourceHandleListener mockSourceHandleListener = 
Mockito.mock(SourceHandleListener.class);
+    final TsBlockSerde mockTsBlockSerde = 
Utils.createMockTsBlockSerde(MOCK_TSBLOCK_SIZE);
+    final IClientManager<TEndPoint, SyncDataNodeMPPDataExchangeServiceClient> 
mockClientManager =
+        Mockito.mock(IClientManager.class);
+    final SyncDataNodeMPPDataExchangeServiceClient mockClient =
+        Mockito.mock(SyncDataNodeMPPDataExchangeServiceClient.class);
+    try {
+      
Mockito.when(mockClientManager.borrowClient(remoteEndpoint)).thenReturn(mockClient);
+      Mockito.doAnswer(
+              invocation -> {
+                final TGetDataBlockRequest request = invocation.getArgument(0);
+                final List<ByteBuffer> shortResponse = new ArrayList<>();
+                for (int i = 0;
+                    i < request.getEndSequenceId() - 
request.getStartSequenceId() - 1;
+                    i++) {
+                  shortResponse.add(ByteBuffer.allocate(0));
+                }
+                return new TGetDataBlockResponse(shortResponse);
+              })
+          .when(mockClient)
+          .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    } catch (ClientManagerException | TException e) {
+      Assert.fail(e.getMessage());
+    }
+
+    final SourceHandle sourceHandle =
+        new SourceHandle(
+            remoteEndpoint,
+            remoteFragmentInstanceId,
+            localFragmentInstanceId,
+            localPlanNodeId,
+            0,
+            mockLocalMemoryManager,
+            Executors.newSingleThreadExecutor(),
+            mockTsBlockSerde,
+            mockSourceHandleListener,
+            mockClientManager);
+    sourceHandle.setRetryIntervalInMs(0L);
+    final Future<?> blocked = sourceHandle.isBlocked();
+
+    sourceHandle.updatePendingDataBlockInfo(
+        0,
+        Stream.generate(() -> MOCK_TSBLOCK_SIZE)
+            .limit(numOfMockTsBlock)
+            .collect(Collectors.toList()));
+
+    try {
+      Mockito.verify(mockClient, 
Mockito.timeout(10_000).times(SourceHandle.MAX_ATTEMPT_TIMES))
+          .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    } catch (TException e) {
+      Assert.fail(e.getMessage());
+    }
+    Mockito.verify(mockSourceHandleListener, Mockito.timeout(10_000).times(1))
+        .onFailure(Mockito.eq(sourceHandle), Mockito.any(TException.class));
+    Assert.assertFalse(blocked.isDone());
+    Assert.assertEquals(0L, sourceHandle.getBufferRetainedSizeInBytes());
+
+    sourceHandle.abort();
+    Assert.assertTrue(blocked.isDone());
+  }
+
   @Test
   public void testForceClose() {
     final String queryId = "q0";
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSenderTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSenderTest.java
new file mode 100644
index 00000000000..a841c35d78f
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/AsyncPlanNodeSenderTest.java
@@ -0,0 +1,81 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.queryengine.plan.scheduler;
+
+import org.apache.iotdb.commons.auth.entity.User;
+import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
+
+import org.junit.Test;
+
+import java.util.Collections;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class AsyncPlanNodeSenderTest {
+
+  @Test
+  public void testOnlyInsertPayloadIsClassifiedAsUserData() {
+    final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
+    final boolean auditLogEnabled = commonConfig.isEnableAuditLog();
+    try {
+      commonConfig.setEnableAuditLog(true);
+
+      final PlanNode queryPlan = mock(PlanNode.class);
+      when(queryPlan.getChildren()).thenReturn(Collections.emptyList());
+      assertFalse(AsyncPlanNodeSender.containsUserData(queryPlan));
+
+      final InsertNode insertNode = mock(InsertNode.class);
+      assertTrue(AsyncPlanNodeSender.containsUserData(insertNode));
+      assertTrue(AsyncPlanNodeSender.containsUserData(insertNode, "root"));
+      assertFalse(
+          AsyncPlanNodeSender.containsUserData(
+              insertNode, User.BUILTIN_INTERNAL_AUDIT_LOG_USERNAME));
+
+      final PlanNode wrapper = mock(PlanNode.class);
+      
when(wrapper.getChildren()).thenReturn(Collections.singletonList(insertNode));
+      assertTrue(AsyncPlanNodeSender.containsUserData(wrapper));
+    } finally {
+      commonConfig.setEnableAuditLog(auditLogEnabled);
+    }
+  }
+
+  @Test
+  public void testAuditDisabledSkipsPlanTraversal() {
+    final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
+    final boolean auditLogEnabled = commonConfig.isEnableAuditLog();
+    final PlanNode planNode = mock(PlanNode.class);
+    try {
+      commonConfig.setEnableAuditLog(false);
+
+      assertFalse(AsyncPlanNodeSender.containsUserData(planNode, "root"));
+      verify(planNode, never()).getChildren();
+    } finally {
+      commonConfig.setEnableAuditLog(auditLogEnabled);
+    }
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/ConsensusReqReaderTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/ConsensusReqReaderTest.java
index e1a48eb73f0..1bd1c2c19f2 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/ConsensusReqReaderTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/ConsensusReqReaderTest.java
@@ -314,6 +314,7 @@ public class ConsensusReqReaderTest {
     Assert.assertEquals(1L, request.getSearchIndex());
     Assert.assertEquals(123456789L, request.getPhysicalTime());
     Assert.assertEquals(7, request.getNodeId());
+    Assert.assertTrue(request.containsUserData());
   }
 
   @Test
@@ -514,6 +515,7 @@ public class ConsensusReqReaderTest {
     PlanNode planNode;
     Assert.assertTrue(iterator.hasNext());
     request = iterator.next();
+    Assert.assertFalse(request.containsUserData());
     Assert.assertEquals(1, request.getRequests().size());
     for (IConsensusRequest innerRequest : request.getRequests()) {
       planNode = 
WALEntry.deserializeForConsensus(innerRequest.serializeToByteBuffer());
diff --git 
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/CommonMessages.java
 
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/CommonMessages.java
index 40c75066d87..bcf90c71e08 100644
--- 
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/CommonMessages.java
+++ 
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/CommonMessages.java
@@ -322,6 +322,10 @@ public final class CommonMessages {
       "Only column with double, float, int32, int64 can be calculated by the 
function, %s is the %s.";
   public static final String 
LOG_TRUSTED_CHANNEL_FUNCTION_FAILED_INITIATOR_ARG_TARGET_ARG_E4C28443 =
       "Trusted channel function failed: initiator=%s, target=%s";
+  public static final String
+      
LOG_USER_DATA_TRANSFER_ATTEMPT_TIME_ARG_INITIATOR_ARG_SOURCE_ARG_TARGET_ARG_PROTECTION_METHOD_ARG_RESULT_ARG_ERROR_ARG_D3E9A1DF
 =
+          "User data transfer attempt: time=%d, initiator=%s, source=%s, 
target=%s,"
+              + " protection_method=%s, result=%s, error=%s";
   public static final String 
EXCEPTION_FILTER_FUNCTION_ROW_INDEX_EXCEED_MAXIMUM =
       "row index exceeds the maximum allowed number in one partition";
   public static final String
diff --git 
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/CommonMessages.java
 
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/CommonMessages.java
index 6cddaa0b1fc..a24b8411ddd 100644
--- 
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/CommonMessages.java
+++ 
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/CommonMessages.java
@@ -217,6 +217,9 @@ public final class CommonMessages {
   public static final String 
EXCEPTION_DISK_SPACE_WARNING_THRESHOLD_MUST_BE_IN_0_1_BUT_WAS_7B345766 = 
"disk_space_warning_threshold 必须在 [0, 1) 范围内,但实际为 ";
   public static final String 
LOG_TRUSTED_CHANNEL_FUNCTION_FAILED_INITIATOR_ARG_TARGET_ARG_E4C28443 =
       "可信信道功能失效:发起者=%s,目标端=%s";
+  public static final String
+      
LOG_USER_DATA_TRANSFER_ATTEMPT_TIME_ARG_INITIATOR_ARG_SOURCE_ARG_TARGET_ARG_PROTECTION_METHOD_ARG_RESULT_ARG_ERROR_ARG_D3E9A1DF
 =
+          "用户数据传送尝试:时间=%d,发起者=%s,源端=%s,目标端=%s,保护方法=%s,结果=%s,错误=%s";
   public static final String EXCEPTION_FILTER_FUNCTION_WPASS_VALIDATION = 
"wpass的取值范围应该位于(0, 1)";
   public static final String EXCEPTION_NO_CALCULATE_COLUMNS = "没有找到可以计算的列.";
   public static final String EXCEPTION_NOT_ALLOWED_COLUMNS = "只允许列类型为double, 
float, int32, int64参与函数计算, 当前列 %s 类型是 %s.";
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AbstractAuditLogger.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AbstractAuditLogger.java
index d4b94265f23..48d0e2811db 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AbstractAuditLogger.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AbstractAuditLogger.java
@@ -35,6 +35,8 @@ public abstract class AbstractAuditLogger {
   private static final long INTERNAL_AUDIT_LOG_USER_ID = 4;
   private static final ThreadLocal<Boolean> RECORDING_TRUSTED_CHANNEL_FAILURE =
       ThreadLocal.withInitial(() -> false);
+  private static final ThreadLocal<Boolean> RECORDING_USER_DATA_TRANSFER =
+      ThreadLocal.withInitial(() -> false);
 
   public static final String OBJECT_AUTHENTICATION_AUDIT_STR =
       "User %s (ID=%d) requests authority on object %s with result %s";
@@ -136,4 +138,51 @@ public abstract class AbstractAuditLogger {
       RECORDING_TRUSTED_CHANNEL_FAILURE.remove();
     }
   }
+
+  /** Records one user-data transfer attempt without retaining any transferred 
payload. */
+  public void recordUserDataTransferAuditLog(UserDataTransferAuditEvent event) 
{
+    if (!IS_AUDIT_LOG_ENABLED
+        || event == null
+        || event.getInitiator() == null
+        || event.getSource() == null
+        || event.getTarget() == null
+        || Boolean.TRUE.equals(RECORDING_USER_DATA_TRANSFER.get())) {
+      return;
+    }
+
+    final String initiatorIdentifier = 
NodeUrlUtils.convertTEndPointUrl(event.getInitiator());
+    final String sourceIdentifier = 
NodeUrlUtils.convertTEndPointUrl(event.getSource());
+    final String targetIdentifier = 
NodeUrlUtils.convertTEndPointUrl(event.getTarget());
+    RECORDING_USER_DATA_TRANSFER.set(true);
+    try {
+      log(
+          createUserDataTransferAuditLogFields(event, initiatorIdentifier),
+          () ->
+              String.format(
+                  CommonMessages
+                      
.LOG_USER_DATA_TRANSFER_ATTEMPT_TIME_ARG_INITIATOR_ARG_SOURCE_ARG_TARGET_ARG_PROTECTION_METHOD_ARG_RESULT_ARG_ERROR_ARG_D3E9A1DF,
+                  event.getTimestamp(),
+                  initiatorIdentifier,
+                  sourceIdentifier,
+                  targetIdentifier,
+                  event.getProtectionMethod(),
+                  event.isSuccess(),
+                  event.getError()));
+    } catch (RuntimeException ignored) {
+      // Audit recording must not affect the user-data transfer being audited.
+    } finally {
+      RECORDING_USER_DATA_TRANSFER.remove();
+    }
+  }
+
+  static AuditLogFields createUserDataTransferAuditLogFields(
+      UserDataTransferAuditEvent event, String initiatorIdentifier) {
+    return new AuditLogFields(
+        INTERNAL_AUDIT_LOG_USER_ID,
+        User.BUILTIN_INTERNAL_AUDIT_LOG_USERNAME,
+        initiatorIdentifier,
+        AuditEventType.USER_DATA_TRANSFER,
+        AuditLogOperation.CONTROL,
+        event.isSuccess());
+  }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
index d9e2a1f5c30..d1f8e63058a 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
@@ -46,6 +46,7 @@ public enum AuditEventType {
   SESSION_TIME_EXCEEDED,
   LOGIN_REJECT_IP,
   TRUSTED_CHANNEL_FUNCTION_FAILURE,
+  USER_DATA_TRANSFER,
   SYSTEM_OPERATION,
 
   DN_SHUTDOWN;
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEvent.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEvent.java
new file mode 100644
index 00000000000..7564ff044d9
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEvent.java
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.commons.audit;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+
+import javax.annotation.Nullable;
+
+/**
+ * Describes one attempt to transfer user data between physically separated 
parts of IoTDB. Payload
+ * contents and exception messages must not be included in this event.
+ */
+public final class UserDataTransferAuditEvent {
+
+  private final long timestamp;
+  private final TEndPoint initiator;
+  private final TEndPoint source;
+  private final TEndPoint target;
+  private final UserDataTransferProtectionMethod protectionMethod;
+  private final boolean success;
+  private final String error;
+
+  public UserDataTransferAuditEvent(
+      TEndPoint initiator,
+      TEndPoint source,
+      TEndPoint target,
+      UserDataTransferProtectionMethod protectionMethod,
+      boolean success,
+      @Nullable String error) {
+    this.timestamp = System.currentTimeMillis();
+    this.initiator = initiator;
+    this.source = source;
+    this.target = target;
+    this.protectionMethod = protectionMethod;
+    this.success = success;
+    this.error = error;
+  }
+
+  public long getTimestamp() {
+    return timestamp;
+  }
+
+  public TEndPoint getInitiator() {
+    return initiator;
+  }
+
+  public TEndPoint getSource() {
+    return source;
+  }
+
+  public TEndPoint getTarget() {
+    return target;
+  }
+
+  public UserDataTransferProtectionMethod getProtectionMethod() {
+    return protectionMethod;
+  }
+
+  public boolean isSuccess() {
+    return success;
+  }
+
+  @Nullable
+  public String getError() {
+    return error;
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditHandler.java
similarity index 54%
copy from 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
copy to 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditHandler.java
index d9e2a1f5c30..171a90cee20 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferAuditHandler.java
@@ -19,39 +19,30 @@
 
 package org.apache.iotdb.commons.audit;
 
-public enum AuditEventType {
-  CHANGE_AUDIT_OPTION,
-  AUDIT_STORAGE_FULL,
-  GENERATE_KEY,
-  DESTROY_KEY,
-  EXECUTE_ENCRYPT,
-  OBJECT_AUTHENTICATION,
-  LBAC_AUTHENTICATION,
-  EXPORT_DATA_WITH_LABEL,
-  IMPORT_DATA_WITH_LABEL,
-  INTEGRITY_CHECK,
-  LOGIN_FAIL_MAX_TIMES,
-  MODIFY_PASSWD,
-  LOGIN,
-  LOGOUT,
-  LOGIN_FINAL,
-  MODIFY_SECURITY_OPTIONS,
-  MODIFY_DEFAULT_SECURITY_VALUES,
-  MODIFY_ROLE_MEMBERSHIP,
-  REVOKE_FAILED,
-  GRANT_ROLE_FAILED,
-  LOGIN_RESOURCE_RESTRICT,
-  LOGIN_FAILED_TRIES,
-  LOGIN_EXCEED_LIMIT,
-  SESSION_TIME_EXCEEDED,
-  LOGIN_REJECT_IP,
-  TRUSTED_CHANNEL_FUNCTION_FAILURE,
-  SYSTEM_OPERATION,
+@FunctionalInterface
+public interface UserDataTransferAuditHandler {
 
-  DN_SHUTDOWN;
+  UserDataTransferAuditHandler NO_OP =
+      new UserDataTransferAuditHandler() {
+        @Override
+        public void onAttempt(UserDataTransferAuditEvent event) {
+          // Do nothing.
+        }
 
-  @Override
-  public String toString() {
-    return name();
+        @Override
+        public boolean isEnabled() {
+          return false;
+        }
+      };
+
+  /**
+   * Records one transfer attempt. Implementations must return promptly and 
must not throw because
+   * callers invoke this method on data-transfer paths.
+   */
+  void onAttempt(UserDataTransferAuditEvent event);
+
+  /** Returns whether transfer audit is enabled. Implementations must not 
throw. */
+  default boolean isEnabled() {
+    return true;
   }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
similarity index 55%
copy from 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
copy to 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
index d9e2a1f5c30..2bf187999ed 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
@@ -19,39 +19,9 @@
 
 package org.apache.iotdb.commons.audit;
 
-public enum AuditEventType {
-  CHANGE_AUDIT_OPTION,
-  AUDIT_STORAGE_FULL,
-  GENERATE_KEY,
-  DESTROY_KEY,
-  EXECUTE_ENCRYPT,
-  OBJECT_AUTHENTICATION,
-  LBAC_AUTHENTICATION,
-  EXPORT_DATA_WITH_LABEL,
-  IMPORT_DATA_WITH_LABEL,
-  INTEGRITY_CHECK,
-  LOGIN_FAIL_MAX_TIMES,
-  MODIFY_PASSWD,
-  LOGIN,
-  LOGOUT,
-  LOGIN_FINAL,
-  MODIFY_SECURITY_OPTIONS,
-  MODIFY_DEFAULT_SECURITY_VALUES,
-  MODIFY_ROLE_MEMBERSHIP,
-  REVOKE_FAILED,
-  GRANT_ROLE_FAILED,
-  LOGIN_RESOURCE_RESTRICT,
-  LOGIN_FAILED_TRIES,
-  LOGIN_EXCEED_LIMIT,
-  SESSION_TIME_EXCEEDED,
-  LOGIN_REJECT_IP,
-  TRUSTED_CHANNEL_FUNCTION_FAILURE,
-  SYSTEM_OPERATION,
-
-  DN_SHUTDOWN;
-
-  @Override
-  public String toString() {
-    return name();
-  }
+public enum UserDataTransferErrorCode {
+  EMPTY_RESPONSE,
+  UNEXPECTED_RESPONSE_SIZE,
+  RECEIVER_CLOSED,
+  REMOTE_REJECTED
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferProtectionMethod.java
similarity index 55%
copy from 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
copy to 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferProtectionMethod.java
index d9e2a1f5c30..9f87afa55b4 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/AuditEventType.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferProtectionMethod.java
@@ -19,39 +19,11 @@
 
 package org.apache.iotdb.commons.audit;
 
-public enum AuditEventType {
-  CHANGE_AUDIT_OPTION,
-  AUDIT_STORAGE_FULL,
-  GENERATE_KEY,
-  DESTROY_KEY,
-  EXECUTE_ENCRYPT,
-  OBJECT_AUTHENTICATION,
-  LBAC_AUTHENTICATION,
-  EXPORT_DATA_WITH_LABEL,
-  IMPORT_DATA_WITH_LABEL,
-  INTEGRITY_CHECK,
-  LOGIN_FAIL_MAX_TIMES,
-  MODIFY_PASSWD,
-  LOGIN,
-  LOGOUT,
-  LOGIN_FINAL,
-  MODIFY_SECURITY_OPTIONS,
-  MODIFY_DEFAULT_SECURITY_VALUES,
-  MODIFY_ROLE_MEMBERSHIP,
-  REVOKE_FAILED,
-  GRANT_ROLE_FAILED,
-  LOGIN_RESOURCE_RESTRICT,
-  LOGIN_FAILED_TRIES,
-  LOGIN_EXCEED_LIMIT,
-  SESSION_TIME_EXCEEDED,
-  LOGIN_REJECT_IP,
-  TRUSTED_CHANNEL_FUNCTION_FAILURE,
-  SYSTEM_OPERATION,
+public enum UserDataTransferProtectionMethod {
+  TLS,
+  NONE;
 
-  DN_SHUTDOWN;
-
-  @Override
-  public String toString() {
-    return name();
+  public static UserDataTransferProtectionMethod fromTlsEnabled(boolean 
tlsEnabled) {
+    return tlsEnabled ? TLS : NONE;
   }
 }
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/AbstractAuditLoggerTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/AbstractAuditLoggerTest.java
index fcb3dbf6d1c..75de1513a9f 100644
--- 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/AbstractAuditLoggerTest.java
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/AbstractAuditLoggerTest.java
@@ -122,6 +122,25 @@ public class AbstractAuditLoggerTest {
     assertSame(auditFailure, channelFailure.getSuppressed()[0]);
   }
 
+  @Test
+  public void testUserDataTransferUsesControlOperation() {
+    final UserDataTransferAuditEvent event =
+        new UserDataTransferAuditEvent(
+            new TEndPoint("127.0.0.1", 10740),
+            new TEndPoint("127.0.0.2", 10740),
+            new TEndPoint("127.0.0.1", 10740),
+            UserDataTransferProtectionMethod.NONE,
+            true,
+            null);
+
+    final AuditLogFields auditLogFields =
+        AbstractAuditLogger.createUserDataTransferAuditLogFields(event, 
"127.0.0.1:10740");
+
+    assertEquals(AuditEventType.USER_DATA_TRANSFER, 
auditLogFields.getAuditEventType());
+    assertEquals(AuditLogOperation.CONTROL, 
auditLogFields.getAuditLogOperation());
+    assertTrue(auditLogFields.getResult());
+  }
+
   private static class TestAuditLogger extends AbstractAuditLogger {
 
     private IAuditEntity auditEntity;
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEventTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEventTest.java
new file mode 100644
index 00000000000..423b05cf9d4
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/audit/UserDataTransferAuditEventTest.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.commons.audit;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+
+import org.junit.Test;
+
+import java.io.IOException;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+
+public class UserDataTransferAuditEventTest {
+
+  @Test
+  public void testRecordsOnlyMinimumTransferFields() {
+    final UserDataTransferAuditEvent event =
+        new UserDataTransferAuditEvent(
+            new TEndPoint("127.0.0.1", 10740),
+            new TEndPoint("127.0.0.2", 10740),
+            new TEndPoint("127.0.0.1", 10740),
+            UserDataTransferProtectionMethod.TLS,
+            false,
+            IOException.class.getName());
+
+    assertEquals(IOException.class.getName(), event.getError());
+    assertFalse(event.isSuccess());
+  }
+
+  @Test
+  public void testProtectionMethodNamesMatchAuditSchema() {
+    assertEquals(
+        UserDataTransferProtectionMethod.TLS,
+        UserDataTransferProtectionMethod.fromTlsEnabled(true));
+    assertEquals(
+        UserDataTransferProtectionMethod.NONE,
+        UserDataTransferProtectionMethod.fromTlsEnabled(false));
+  }
+}

Reply via email to