| | |
| | | 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 |
| | |
| | | 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()); |
| | |
| | | dto.setT(ptz.getT()); |
| | | dto.setZ(ptz.getZ()); |
| | | } |
| | | |
| | | return dto; |
| | | } |
| | | |
| | |
| | | && !EXCLUDE_FACTORY.contains(camera.getFactory()) |
| | | && SUPPORT_TYPES.contains(camera.getType()); |
| | | } |
| | | } |
| | | } |