Prechádzať zdrojové kódy

[0924] 文件扫描任务兼容已存在数据,批量处理所有范围内的数据

taosj 6 dní pred
rodič
commit
8145351491

+ 8 - 0
forecast-backend-server/src/main/java/com/ufo/project/data/mapper/DataFileSyncFilesMapper.java

@@ -19,6 +19,14 @@ public interface DataFileSyncFilesMapper
      */
     public DataFileSyncFiles selectDataFileSyncFilesById(Long id);
 
+    /**
+     * 按文件路径(精确匹配)查询待处理数据文件同步
+     *
+     * @param file 文件(完整路径)
+     * @return 待处理数据文件同步, 不存在返回null
+     */
+    public DataFileSyncFiles selectDataFileSyncFilesByFile(String file);
+
     /**
      * 查询待处理数据文件同步列表
      *

+ 8 - 0
forecast-backend-server/src/main/java/com/ufo/project/data/service/IDataFileSyncFilesService.java

@@ -19,6 +19,14 @@ public interface IDataFileSyncFilesService
      */
     public DataFileSyncFiles selectDataFileSyncFilesById(Long id);
 
+    /**
+     * 按文件路径(精确匹配)查询待处理数据文件同步
+     *
+     * @param file 文件(完整路径)
+     * @return 待处理数据文件同步, 不存在返回null
+     */
+    public DataFileSyncFiles selectDataFileSyncFilesByFile(String file);
+
     /**
      * 查询待处理数据文件同步列表
      *

+ 12 - 0
forecast-backend-server/src/main/java/com/ufo/project/data/service/impl/DataFileSyncFilesServiceImpl.java

@@ -33,6 +33,18 @@ public class DataFileSyncFilesServiceImpl implements IDataFileSyncFilesService
         return dataFileSyncFilesMapper.selectDataFileSyncFilesById(id);
     }
 
+    /**
+     * 按文件路径(精确匹配)查询待处理数据文件同步
+     *
+     * @param file 文件(完整路径)
+     * @return 待处理数据文件同步, 不存在返回null
+     */
+    @Override
+    public DataFileSyncFiles selectDataFileSyncFilesByFile(String file)
+    {
+        return dataFileSyncFilesMapper.selectDataFileSyncFilesByFile(file);
+    }
+
     /**
      * 查询待处理数据文件同步列表
      *

+ 15 - 13
forecast-backend-server/src/main/java/com/ufo/project/data/task/DataFileSyncTask.java

@@ -6,7 +6,6 @@ import java.nio.file.Path;
 import java.nio.file.Paths;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicInteger;
 import java.util.stream.Stream;
 
 import com.alibaba.fastjson2.JSON;
@@ -124,28 +123,22 @@ public class DataFileSyncTask
         }
         long intervalHours = config.getUpdateInterval() == null ? 0L : config.getUpdateInterval();
         long modifiedAfter = System.currentTimeMillis() - TimeUnit.HOURS.toMillis(intervalHours);
-        AtomicInteger inserted = new AtomicInteger();
+
         try (Stream<Path> stream = Files.walk(root))
         {
             stream.filter(Files::isRegularFile)
                   .filter(path -> path.toFile().lastModified() >= modifiedAfter)
                   .forEach(path -> {
-                      if (registerFile(config, path))
-                      {
-                          inserted.incrementAndGet();
-                      }
+                      registerFile(config, path);
                   });
         }
-        if (inserted.get() > 0)
-        {
-            log.info("[DataFileSyncTask] 目录 {} 新登记待处理文件 {} 个", config.getScanFolder(), inserted.get());
-        }
     }
 
     /**
-     * 登记单个文件, 已存在(按文件路径查重)则跳过
+     * 登记单个文件: 先按文件路径精确查重, 已存在相同file记录则直接返回true;
+     * 不存在则保存到data_file_sync_files(保留insert-if-absent作为并发兜底)
      *
-     * @return 是否实际插入
+     * @return 文件是否已登记(已存在或本次插入成功)
      */
     private boolean registerFile(FileSyncFolderConfig config, Path path)
     {
@@ -158,7 +151,16 @@ public class DataFileSyncTask
             row.setHandler(config.getHandler());
             row.setProcessTimes(0);
             row.setPriority(DataFileSyncFiles.DEFAULT_PRIORITY);
-            return dataFileSyncFilesService.insertDataFileSyncFilesIfAbsent(row) > 0;
+            if (dataFileSyncFilesService.selectDataFileSyncFilesByFile(row.getFile()) != null)
+            {
+                return true;
+            }
+            if (dataFileSyncFilesService.insertDataFileSyncFilesIfAbsent(row) > 0)
+            {
+                log.info("[DataFileSyncTask] 新登记待处理文件, file={}, handler={}", row.getFile(), row.getHandler());
+                return true;
+            }
+            return false;
         }
         catch (Exception e)
         {

+ 6 - 0
forecast-backend-server/src/main/resources/mybatis/data/DataFileSyncFilesMapper.xml

@@ -43,6 +43,12 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
         where id = #{id}
     </select>
 
+    <select id="selectDataFileSyncFilesByFile" parameterType="String" resultMap="DataFileSyncFilesResult">
+        <include refid="selectDataFileSyncFilesVo"/>
+        where file = #{file}
+        limit 1
+    </select>
+
     <insert id="insertDataFileSyncFilesIfAbsent" parameterType="DataFileSyncFiles">
         insert into data_file_sync_files
             (type, file, status, handler, create_time, update_time, process_times, priority)