liusuyi
2026-05-30 f4f4fc53260eb67483dce406a963628273786a61
ard-modules/ard-modules-work/src/main/java/com/ard/work/device/camera/service/PtzDataCollector.java
@@ -12,10 +12,14 @@
import org.springframework.scheduling.annotation.Async;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Component;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
@Slf4j
@Component
@@ -25,37 +29,69 @@
    private final CameraSDKService cameraSDKService;
    private final PtzCacheManager ptzCacheManager;
    @Resource(name = "timeoutExecutor")
    private ThreadPoolTaskExecutor sdkExecutor;
    private static final Set<String> EXCLUDE_FACTORY = new HashSet<>(Arrays.asList("7"));
    private static final String CAMERA_ENABLE = "1";
    private static final Set<String> SUPPORT_TYPES = new HashSet<>(Arrays.asList("1", "4"));
    /** 单通道 SDK 调用超时 */
    private static final long SDK_CALL_TIMEOUT_MS = 2000L;
    /** 整相机所有通道总采集超时 */
    private static final long CAMERA_TOTAL_TIMEOUT_MS = 3000L;
    @Async("cameraPTZExecutor")
    public void collectCamera(ArdCamera camera, Set<String> runningCamera) {
        String camId = camera.getId();
        try {
            if (!isValidCamera(camera)) return;
            List<ArdChannel> channels = camera.getChannelList();
            if (channels == null || channels.isEmpty()) return;
            for (ArdChannel channel : channels) {
                String key = camId + "_" + channel.getChanNo();
                PtzParamDTO dto= fetchPtz(camera, channel);
                if (dto != null) {
                    ptzCacheManager.updatePtzData(key, dto);
                }
            CompletableFuture<?>[] futures = channels.stream()
                    .map(ch -> CompletableFuture.runAsync(
                            () -> collectChannel(camera, ch), sdkExecutor))
                    .toArray(CompletableFuture[]::new);
            try {
                CompletableFuture.allOf(futures).get(CAMERA_TOTAL_TIMEOUT_MS, TimeUnit.MILLISECONDS);
            } catch (TimeoutException e) {
                log.warn("采集相机 {} 整体超时,部分通道数据可能未更新", camId);
            } catch (Exception e) {
                log.error("采集相机 {} 异常", camId, e);
            }
        } finally {
            // ✅ 必须释放(防死锁)
            runningCamera.remove(camId);
        }
    }
    /**
     * 原始SDK调用(不改)
     */
    private PtzParamDTO fetchPtz(ArdCamera camera, ArdChannel channel) {
    private void collectChannel(ArdCamera camera, ArdChannel channel) {
        PtzParamDTO dto = fetchPtzWithTimeout(camera, channel);
        if (dto != null) {
            ptzCacheManager.updatePtzData(camera.getId() + "_" + channel.getChanNo(), dto);
        }
    }
    /**
     * 带超时的 SDK 调用,防止 SDK 卡死拖垮线程池
     */
    private PtzParamDTO fetchPtzWithTimeout(ArdCamera camera, ArdChannel channel) {
        CompletableFuture<PtzParamDTO> future = CompletableFuture.supplyAsync(
                () -> fetchPtz(camera, channel), sdkExecutor);
        try {
            return future.get(SDK_CALL_TIMEOUT_MS, TimeUnit.MILLISECONDS);
        } catch (TimeoutException e) {
            log.warn("相机 {} 通道 {} SDK 调用超时", camera.getId(), channel.getChanNo());
            future.cancel(true);
            return null;
        } catch (Exception e) {
            log.error("相机 {} 通道 {} SDK 调用异常", camera.getId(), channel.getChanNo(), e);
            return null;
        }
    }
    private PtzParamDTO fetchPtz(ArdCamera camera, ArdChannel channel) {
        CameraCmd cmd = new CameraCmd();
        cmd.setCameraId(camera.getId());
        cmd.setChanNo(channel.getChanNo());
@@ -78,7 +114,6 @@
            dto.setT(ptz.getT());
            dto.setZ(ptz.getZ());
        }
        return dto;
    }
@@ -87,4 +122,4 @@
                && !EXCLUDE_FACTORY.contains(camera.getFactory())
                && SUPPORT_TYPES.contains(camera.getType());
    }
}
}