修复大华车辆通行记录增量同步

dev
shenzhidan 2026-08-06 15:11:56 +08:00
parent f281811fe6
commit c7ebddef50
28 changed files with 1447 additions and 3 deletions

View File

@ -0,0 +1,434 @@
package com.zcloud.primeport.service;
import cn.hutool.json.JSONUtil;
import com.zcloud.primeport.domain.gateway.DaHuaGateway;
import com.zcloud.primeport.domain.gateway.DaHuaResourceRepositoryGateway;
import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway;
import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordSyncGateway;
import com.zcloud.primeport.domain.model.CorpInfoSnapshotE;
import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE;
import com.zcloud.primeport.domain.model.VehicleApplySnapshotE;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.time.format.DateTimeParseException;
import java.text.Normalizer;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.function.Consumer;
@Service
@RequiredArgsConstructor
public class DaHuaVehicleAccessRecordSyncApplicationService
implements DaHuaVehicleAccessRecordSyncGateway {
private static final DateTimeFormatter DATE_TIME_FORMATTER =
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
private static final int MAX_EXCEPTION_DETAILS = 100;
private static final int MAX_PAGE_COUNT = 10000;
private final DaHuaGateway daHuaGateway;
private final DaHuaResourceRepositoryGateway resourceRepositoryGateway;
private final DaHuaVehicleAccessRecordRepositoryGateway recordRepositoryGateway;
@Override
public DaHuaVehicleAccessRecordSyncResultE syncIncremental(int pageSize, int delayMinutes,
int overlapMinutes,
int initialLookbackMinutes) {
validatePageSize(pageSize);
LocalDateTime end = LocalDateTime.now().minusMinutes(Math.max(delayMinutes, 0));
LocalDateTime checkpoint = recordRepositoryGateway.findSyncCheckpoint();
LocalDateTime start = checkpoint == null
? end.minusMinutes(Math.max(initialLookbackMinutes, 1))
: checkpoint.minusMinutes(Math.max(overlapMinutes, 0));
if (!start.isBefore(end)) {
return new DaHuaVehicleAccessRecordSyncResultE();
}
DaHuaVehicleAccessRecordSyncResultE result = sync(format(start), format(end), pageSize);
recordRepositoryGateway.saveSyncCheckpoint(end);
return result;
}
@Override
public DaHuaVehicleAccessRecordSyncResultE sync(String startTime, String endTime, int pageSize) {
LocalDateTime start = parseRequired(startTime, "startTime");
LocalDateTime end = parseRequired(endTime, "endTime");
if (start.isAfter(end)) {
throw new IllegalArgumentException("startTime cannot be after endTime");
}
validatePageSize(pageSize);
SyncIndex index = buildIndex();
Map<String, DaHuaVehicleAccessRecordE> synchronizedRecords = new HashMap<>();
DaHuaVehicleAccessRecordSyncResultE result = new DaHuaVehicleAccessRecordSyncResultE();
syncWindow("enterTimeStrLeft", "enterTimeStrRight", format(start), format(end),
pageSize, index, synchronizedRecords, result);
syncWindow("exitTimeStrLeft", "exitTimeStrRight", format(start), format(end),
pageSize, index, synchronizedRecords, result);
if (result.getRecordInvalid() > 0) {
throw new IllegalStateException("Dahua vehicle access record sync contains "
+ result.getRecordInvalid() + " invalid record(s): "
+ String.join("; ", result.getExceptionDetails()));
}
return result;
}
private void syncWindow(String leftField, String rightField, String start, String end,
int pageSize, SyncIndex index,
Map<String, DaHuaVehicleAccessRecordE> synchronizedRecords,
DaHuaVehicleAccessRecordSyncResultE result) {
Set<String> seenPageSignatures = new HashSet<>();
for (int pageNum = 1; pageNum <= MAX_PAGE_COUNT; pageNum++) {
Map<String, Object> params = new LinkedHashMap<>();
params.put("pageNum", pageNum);
params.put("pageSize", pageSize);
params.put(leftField, start);
params.put(rightField, end);
Map<String, Object> response = daHuaGateway.queryVehicleAccessRecord(params);
if (response == null || !isSuccessful(response)) {
throw new IllegalStateException("Dahua vehicle access query failed: code="
+ (response == null ? null : response.get("code")) + ", errMsg="
+ (response == null ? "empty response" : response.get("errMsg")));
}
Map<String, Object> data = toMap(response.get("data"));
List<Map<String, Object>> pageData = toRecordList(data.get("pageData"));
result.setPageCount(result.getPageCount() + 1);
if (pageData.isEmpty()) {
return;
}
String pageSignature = JSONUtil.toJsonStr(pageData);
if (!seenPageSignatures.add(pageSignature)) {
throw new IllegalStateException("Dahua vehicle access query returned a repeated page: "
+ pageNum);
}
persistPage(pageData, index, synchronizedRecords, result);
}
throw new IllegalStateException("Dahua vehicle access query exceeded maximum page count: "
+ MAX_PAGE_COUNT);
}
private void persistPage(List<Map<String, Object>> pageData, SyncIndex index,
Map<String, DaHuaVehicleAccessRecordE> synchronizedRecords,
DaHuaVehicleAccessRecordSyncResultE result) {
Set<String> accessIds = new LinkedHashSet<>();
Set<String> plates = new LinkedHashSet<>();
for (Map<String, Object> source : pageData) {
String accessId = blankToNull(firstString(source, "accessId", "id"));
if (accessId != null && !synchronizedRecords.containsKey(accessId)) {
accessIds.add(accessId);
}
addPlate(plates, firstString(source, "carNum"));
addPlate(plates, firstString(source, "exitCarNum"));
}
for (DaHuaVehicleAccessRecordE existing : recordRepositoryGateway.listByDahuaAccessIds(accessIds)) {
synchronizedRecords.put(existing.getDahuaAccessId(), existing);
}
for (VehicleApplySnapshotE snapshot : nullSafe(resourceRepositoryGateway.listVehicleApplySnapshots(plates))) {
index.vehicleSnapshotsByPlate.computeIfAbsent(normalizePlate(snapshot.getLicenceNo()),
ignored -> new ArrayList<>()).add(snapshot);
}
for (Map<String, Object> source : pageData) {
result.setRecordTotal(result.getRecordTotal() + 1);
try {
String accessId = blankToNull(firstString(source, "accessId", "id"));
if (accessId == null) {
throw new IllegalArgumentException("accessId is missing");
}
DaHuaVehicleAccessRecordE record = synchronizedRecords.get(accessId);
boolean created = record == null;
if (created) {
record = new DaHuaVehicleAccessRecordE();
}
mapRecord(record, source, accessId, index);
if (created) {
recordRepositoryGateway.add(record);
result.setRecordCreated(result.getRecordCreated() + 1);
} else {
recordRepositoryGateway.update(record);
result.setRecordUpdated(result.getRecordUpdated() + 1);
}
synchronizedRecords.put(accessId, record);
} catch (Exception e) {
result.setRecordInvalid(result.getRecordInvalid() + 1);
if (result.getExceptionDetails().size() < MAX_EXCEPTION_DETAILS) {
result.getExceptionDetails().add("Vehicle access record sync failed: " + e.getMessage());
}
}
}
}
private void mapRecord(DaHuaVehicleAccessRecordE record, Map<String, Object> source,
String accessId, SyncIndex index) {
String departmentName = blankToNull(firstString(source, "departmentName", "name"));
Long departmentId = firstLong(source, "departmentId");
String orgCode = blankToNull(firstString(source, "orgCode"));
DaHuaDepartmentCorpMappingE mapping = resolveMapping(index, departmentId, orgCode, departmentName);
VehicleApplySnapshotE vehicle = resolveVehicle(index, source);
CorpInfoSnapshotE corp = mapping == null || mapping.getCorpId() == null
? null : index.corpsById.get(mapping.getCorpId());
if (mapping != null) {
setIfNotNull(record::setCorpId, mapping.getCorpId());
setIfNotNull(record::setCorpName,
blankToNull(corp == null ? mapping.getCorpName() : corp.getCorpName()));
setIfNotNull(record::setPortArea, corp == null ? null : corp.getPortArea());
} else if (vehicle != null) {
corp = vehicle.getVehicleCorpId() == null ? null : index.corpsById.get(vehicle.getVehicleCorpId());
setIfNotNull(record::setCorpId, vehicle.getVehicleCorpId());
setIfNotNull(record::setCorpName,
blankToNull(corp == null ? vehicle.getVehicleCorpName() : corp.getCorpName()));
setIfNotNull(record::setPortArea, corp == null ? null : corp.getPortArea());
}
record.setDahuaAccessId(accessId);
setIfNotNull(record::setDahuaDepartmentId, departmentId != null ? departmentId
: mapping == null ? (vehicle == null ? null : vehicle.getVehicleDepartmentId())
: mapping.getDahuaDepartmentId());
setIfNotNull(record::setDahuaDepartmentName, departmentName != null ? departmentName
: blankToNull(mapping == null ? (vehicle == null ? null : vehicle.getVehicleDepartmentName())
: mapping.getDahuaDepartmentName()));
setIfNotNull(record::setDahuaOrgCode,
orgCode != null ? orgCode : blankToNull(mapping == null ? null : mapping.getDahuaOrgCode()));
setIfNotNull(record::setParkingLotCode, blankToNull(firstString(source, "parkingLotCode")));
setIfNotNull(record::setParkingLotName, blankToNull(firstString(source, "parkingLot")));
setIfNotNull(record::setCarType, firstInteger(source, "carType"));
setIfNotNull(record::setCarTypeName, blankToNull(firstString(source, "carTypeStr")));
setIfNotNull(record::setCarNum, blankToNull(firstString(source, "carNum")));
setIfNotNull(record::setExitCarNum, blankToNull(firstString(source, "exitCarNum")));
setIfNotNull(record::setEnterItcChannelCode, blankToNull(firstString(source, "enterItcDevChnid")));
setIfNotNull(record::setEnterSluiceChannelCode, blankToNull(firstString(source, "enterSluiceDevChnid")));
setIfNotNull(record::setEnterSluiceChannelName, blankToNull(firstString(source, "enterSluiceDevChnname")));
setIfNotNull(record::setEnterTime, firstDateTime(source, "enterTimeStr", "enterTime"));
setIfNotNull(record::setEnterImageUrl, blankToNull(firstString(source, "enterImg")));
setIfNotNull(record::setExitItcChannelCode, blankToNull(firstString(source, "exitItcDevChnid")));
setIfNotNull(record::setExitSluiceChannelCode, blankToNull(firstString(source, "exitSluiceDevChnid")));
setIfNotNull(record::setExitSluiceChannelName, blankToNull(firstString(source, "exitSluiceDevChnname")));
setIfNotNull(record::setExitTime, firstDateTime(source, "exitTimeStr", "exitTime"));
setIfNotNull(record::setExitImageUrl, blankToNull(firstString(source, "exitImg")));
record.setSourceType("PULL");
record.setRawData(JSONUtil.toJsonStr(source));
record.setLastSyncTime(LocalDateTime.now());
record.setDeleteEnum("FALSE");
}
private SyncIndex buildIndex() {
SyncIndex index = new SyncIndex();
for (CorpInfoSnapshotE corp : nullSafe(resourceRepositoryGateway.listActiveCorps())) {
index.corpsById.put(corp.getId(), corp);
}
for (DaHuaDepartmentCorpMappingE mapping : nullSafe(resourceRepositoryGateway.listMappings())) {
if (mapping == null || "TRUE".equalsIgnoreCase(mapping.getDeleteEnum())
|| (mapping.getBindStatus() != null && mapping.getBindStatus() != 1)) {
continue;
}
if (mapping.getDahuaDepartmentId() != null) {
index.mappingsByDepartmentId.put(mapping.getDahuaDepartmentId(), mapping);
}
if (blankToNull(mapping.getDahuaOrgCode()) != null) {
index.mappingsByOrgCode.put(mapping.getDahuaOrgCode(), mapping);
}
String name = normalizeName(mapping.getDahuaDepartmentName());
if (name != null) {
index.mappingsByName.computeIfAbsent(name, ignored -> new ArrayList<>()).add(mapping);
}
}
return index;
}
private DaHuaDepartmentCorpMappingE resolveMapping(SyncIndex index, Long departmentId,
String orgCode, String departmentName) {
DaHuaDepartmentCorpMappingE mapping = departmentId == null ? null
: index.mappingsByDepartmentId.get(departmentId);
if (mapping == null && orgCode != null) {
mapping = index.mappingsByOrgCode.get(orgCode);
}
if (mapping == null) {
List<DaHuaDepartmentCorpMappingE> matches = index.mappingsByName.get(normalizeName(departmentName));
mapping = matches != null && matches.size() == 1 ? matches.get(0) : null;
}
return mapping;
}
private VehicleApplySnapshotE resolveVehicle(SyncIndex index, Map<String, Object> source) {
VehicleApplySnapshotE enter = uniqueVehicle(index.vehicleSnapshotsByPlate,
firstString(source, "carNum"));
VehicleApplySnapshotE exit = uniqueVehicle(index.vehicleSnapshotsByPlate,
firstString(source, "exitCarNum"));
return enter != null ? enter : exit;
}
private VehicleApplySnapshotE uniqueVehicle(Map<String, List<VehicleApplySnapshotE>> snapshots,
String licenceNo) {
List<VehicleApplySnapshotE> matches = snapshots.get(normalizePlate(licenceNo));
if (matches == null || matches.isEmpty()) {
return null;
}
Set<Long> corpIds = new HashSet<>();
for (VehicleApplySnapshotE match : matches) {
if (match.getVehicleCorpId() != null) {
corpIds.add(match.getVehicleCorpId());
}
}
return corpIds.size() <= 1 ? matches.get(0) : null;
}
private void addPlate(Set<String> plates, String plate) {
if (blankToNull(plate) != null) {
plates.add(plate);
String normalized = normalizePlate(plate);
if (normalized != null) {
plates.add(normalized);
}
}
}
private boolean isSuccessful(Map<String, Object> response) {
if (Boolean.TRUE.equals(response.get("passed")) || Boolean.TRUE.equals(response.get("success"))) {
return true;
}
String code = String.valueOf(response.get("code"));
return "0".equals(code) || "200".equals(code);
}
private LocalDateTime firstDateTime(Map<String, Object> source, String textField, String epochField) {
Object value = source.get(textField);
if (value == null) {
value = source.get(epochField);
}
if (value == null) {
return null;
}
String text = blankToNull(value.toString());
if (text == null) {
return null;
}
if (value instanceof Number || text.matches("\\d{10,}")) {
long epoch = Long.parseLong(text);
Instant instant = text.length() <= 10
? Instant.ofEpochSecond(epoch) : Instant.ofEpochMilli(epoch);
return instant
.atZone(ZoneId.systemDefault()).toLocalDateTime();
}
try {
return LocalDateTime.parse(text, DATE_TIME_FORMATTER);
} catch (DateTimeParseException e) {
throw new IllegalArgumentException(textField + " must use yyyy-MM-dd HH:mm:ss", e);
}
}
private LocalDateTime parseRequired(String value, String field) {
if (blankToNull(value) == null) {
throw new IllegalArgumentException(field + " cannot be blank");
}
try {
return LocalDateTime.parse(value.trim(), DATE_TIME_FORMATTER);
} catch (DateTimeParseException e) {
throw new IllegalArgumentException(field + " must use yyyy-MM-dd HH:mm:ss", e);
}
}
private String format(LocalDateTime value) {
return value.format(DATE_TIME_FORMATTER);
}
private void validatePageSize(int pageSize) {
if (pageSize < 1 || pageSize > 1000) {
throw new IllegalArgumentException("pageSize must be between 1 and 1000");
}
}
private Object firstValue(Map<String, Object> source, String... fields) {
for (String field : fields) {
if (source.get(field) != null) {
return source.get(field);
}
}
return null;
}
private String firstString(Map<String, Object> source, String... fields) {
Object value = firstValue(source, fields);
return value == null ? null : value.toString();
}
private Long firstLong(Map<String, Object> source, String... fields) {
Object value = firstValue(source, fields);
return value == null ? null : Long.valueOf(value.toString());
}
private Integer firstInteger(Map<String, Object> source, String... fields) {
Object value = firstValue(source, fields);
return value == null ? null : Integer.valueOf(value.toString());
}
@SuppressWarnings("unchecked")
private Map<String, Object> toMap(Object value) {
if (value instanceof Map) {
return new LinkedHashMap<>((Map<String, Object>) value);
}
return value == null ? new LinkedHashMap<>() : JSONUtil.parseObj(JSONUtil.toJsonStr(value));
}
@SuppressWarnings("unchecked")
private List<Map<String, Object>> toRecordList(Object value) {
if (!(value instanceof Iterable)) {
return Collections.emptyList();
}
List<Map<String, Object>> result = new ArrayList<>();
for (Object item : (Iterable<?>) value) {
result.add(item instanceof Map ? new LinkedHashMap<>((Map<String, Object>) item)
: JSONUtil.parseObj(JSONUtil.toJsonStr(item)));
}
return result;
}
private <T> Collection<T> nullSafe(Collection<T> values) {
return values == null ? Collections.emptyList() : values;
}
private <T> void setIfNotNull(Consumer<T> setter, T value) {
if (value != null) {
setter.accept(value);
}
}
private String normalizeName(String value) {
String normalized = blankToNull(value);
return normalized == null ? null : Normalizer.normalize(normalized, Normalizer.Form.NFKC)
.replaceAll("\\s+", "").toLowerCase(Locale.ROOT);
}
private String normalizePlate(String value) {
String normalized = blankToNull(value);
return normalized == null ? null : normalized.replaceAll("\\s+", "").toUpperCase(Locale.ROOT);
}
private String blankToNull(String value) {
return value == null || value.trim().isEmpty() ? null : value.trim();
}
private static class SyncIndex {
private final Map<Long, CorpInfoSnapshotE> corpsById = new HashMap<>();
private final Map<Long, DaHuaDepartmentCorpMappingE> mappingsByDepartmentId = new HashMap<>();
private final Map<String, DaHuaDepartmentCorpMappingE> mappingsByOrgCode = new HashMap<>();
private final Map<String, List<DaHuaDepartmentCorpMappingE>> mappingsByName = new HashMap<>();
private final Map<String, List<VehicleApplySnapshotE>> vehicleSnapshotsByPlate = new HashMap<>();
}
}

View File

@ -0,0 +1,319 @@
package com.zcloud.primeport.service;
import com.zcloud.primeport.domain.gateway.DaHuaGateway;
import com.zcloud.primeport.domain.gateway.DaHuaResourceRepositoryGateway;
import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway;
import com.zcloud.primeport.domain.model.CorpInfoSnapshotE;
import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE;
import com.zcloud.primeport.domain.model.VehicleApplySnapshotE;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.anyMap;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class DaHuaVehicleAccessRecordSyncApplicationServiceTest {
@Mock
private DaHuaGateway daHuaGateway;
@Mock
private DaHuaResourceRepositoryGateway resourceRepositoryGateway;
@Mock
private DaHuaVehicleAccessRecordRepositoryGateway recordRepositoryGateway;
private DaHuaVehicleAccessRecordSyncApplicationService applicationService;
@BeforeEach
void setUp() {
applicationService = new DaHuaVehicleAccessRecordSyncApplicationService(
daHuaGateway, resourceRepositoryGateway, recordRepositoryGateway);
when(resourceRepositoryGateway.listActiveCorps()).thenReturn(Collections.emptyList());
when(resourceRepositoryGateway.listMappings()).thenReturn(Collections.emptyList());
lenient().when(resourceRepositoryGateway.listVehicleApplySnapshots(anyCollection()))
.thenReturn(Collections.emptyList());
lenient().when(recordRepositoryGateway.listByDahuaAccessIds(anyCollection()))
.thenReturn(Collections.emptyList());
}
@Test
void shouldQueryBothWindowsAndUpdateTheSameAccessIdIdempotently() {
Map<String, Object> enter = record("A-1");
enter.put("departmentId", 10L);
enter.put("enterSluiceDevChnid", "CH-ENTER");
enter.put("enterTimeStr", "2026-08-01 10:00:00");
Map<String, Object> exit = record("A-1");
exit.put("departmentId", 10L);
exit.put("exitSluiceDevChnid", "CH-EXIT");
exit.put("exitTimeStr", "2026-08-01 18:00:00");
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
if (Integer.valueOf(1).equals(params.get("pageNum"))
&& params.containsKey("enterTimeStrLeft")) {
return response(Collections.singletonList(enter));
}
if (Integer.valueOf(1).equals(params.get("pageNum"))
&& params.containsKey("exitTimeStrLeft")) {
return response(Collections.singletonList(exit));
}
return response(Collections.emptyList());
});
DaHuaVehicleAccessRecordSyncResultE result = applicationService.sync(
"2026-08-01 00:00:00", "2026-08-02 00:00:00", 100);
assertEquals(4, result.getPageCount());
assertEquals(2, result.getRecordTotal());
assertEquals(1, result.getRecordCreated());
assertEquals(1, result.getRecordUpdated());
verify(recordRepositoryGateway).add(any(DaHuaVehicleAccessRecordE.class));
ArgumentCaptor<DaHuaVehicleAccessRecordE> updateCaptor =
ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class);
verify(recordRepositoryGateway).update(updateCaptor.capture());
assertEquals("A-1", updateCaptor.getValue().getDahuaAccessId());
assertEquals("2026-08-01 18:00:00",
updateCaptor.getValue().getExitTime().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
assertEquals("2026-08-01 10:00:00",
updateCaptor.getValue().getEnterTime().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
assertEquals("CH-ENTER", updateCaptor.getValue().getEnterSluiceChannelCode());
verify(daHuaGateway, times(4)).queryVehicleAccessRecord(anyMap());
}
@Test
void shouldFetchEveryPageForEachTimeWindow() {
Map<String, Object> first = record("A-1");
Map<String, Object> second = record("A-2");
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
Integer pageNum = (Integer) params.get("pageNum");
if (params.containsKey("enterTimeStrLeft")) {
return pageNum == 1 ? response(Collections.singletonList(first))
: pageNum == 2 ? response(Collections.singletonList(second))
: response(Collections.emptyList());
}
return response(Collections.emptyList());
});
DaHuaVehicleAccessRecordSyncResultE result = applicationService.sync(
"2026-08-01 00:00:00", "2026-08-02 00:00:00", 1);
assertEquals(4, result.getPageCount());
assertEquals(2, result.getRecordTotal());
assertEquals(2, result.getRecordCreated());
verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 1));
verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 2));
verify(daHuaGateway).queryVehicleAccessRecord(params("enterTimeStrLeft", 3));
verify(daHuaGateway).queryVehicleAccessRecord(params("exitTimeStrLeft", 1));
}
@Test
void shouldResolveCorpByDepartmentAndKeepChannelData() {
CorpInfoSnapshotE personCorp = corp(100L, "人员企业", 1);
when(resourceRepositoryGateway.listActiveCorps())
.thenReturn(Collections.singletonList(personCorp));
DaHuaDepartmentCorpMappingE mapping = new DaHuaDepartmentCorpMappingE();
mapping.setDahuaDepartmentId(10L);
mapping.setDahuaDepartmentName("部门");
mapping.setCorpId(100L);
mapping.setCorpName("人员企业");
mapping.setBindStatus(1);
mapping.setDeleteEnum("FALSE");
when(resourceRepositoryGateway.listMappings()).thenReturn(Collections.singletonList(mapping));
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
return Integer.valueOf(1).equals(params.get("pageNum"))
&& params.containsKey("enterTimeStrLeft")
? response(Collections.singletonList(recordWithDepartmentAndChannel("A-1")))
: response(Collections.emptyList());
});
applicationService.sync("2026-08-01 00:00:00", "2026-08-02 00:00:00", 100);
ArgumentCaptor<DaHuaVehicleAccessRecordE> captor =
ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class);
verify(recordRepositoryGateway).add(captor.capture());
DaHuaVehicleAccessRecordE saved = captor.getValue();
assertEquals(Long.valueOf(100L), saved.getCorpId());
assertEquals("人员企业", saved.getCorpName());
assertEquals(Integer.valueOf(1), saved.getPortArea());
assertEquals("CH-ENTER", saved.getEnterSluiceChannelCode());
}
@Test
void shouldResolveCorpByVehiclePlateWhenDepartmentIsUnmapped() {
VehicleApplySnapshotE vehicle = new VehicleApplySnapshotE();
vehicle.setLicenceNo("冀A12345");
vehicle.setVehicleCorpId(300L);
vehicle.setVehicleCorpName("车辆企业");
vehicle.setVehicleDepartmentId(30L);
vehicle.setVehicleDepartmentName("车辆部门");
when(resourceRepositoryGateway.listVehicleApplySnapshots(anyCollection()))
.thenReturn(Collections.singletonList(vehicle));
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
return Integer.valueOf(1).equals(params.get("pageNum"))
&& params.containsKey("enterTimeStrLeft")
? response(Collections.singletonList(recordWithPlate("A-1", "冀A12345")))
: response(Collections.emptyList());
});
applicationService.sync("2026-08-01 00:00:00", "2026-08-02 00:00:00", 100);
ArgumentCaptor<DaHuaVehicleAccessRecordE> captor =
ArgumentCaptor.forClass(DaHuaVehicleAccessRecordE.class);
verify(recordRepositoryGateway).add(captor.capture());
assertEquals(Long.valueOf(300L), captor.getValue().getCorpId());
assertEquals("车辆企业", captor.getValue().getCorpName());
assertEquals(Long.valueOf(30L), captor.getValue().getDahuaDepartmentId());
assertEquals("车辆部门", captor.getValue().getDahuaDepartmentName());
}
@Test
void shouldFailForMissingAccessIdAndRejectOversizedPage() {
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
return params.containsKey("enterTimeStrLeft")
&& Integer.valueOf(1).equals(params.get("pageNum"))
? response(Collections.singletonList(record(null)))
: response(Collections.emptyList());
});
IllegalStateException exception = assertThrows(IllegalStateException.class,
() -> applicationService.sync(
"2026-08-01 00:00:00", "2026-08-02 00:00:00", 100));
assertTrue(exception.getMessage().contains("accessId is missing"));
verify(recordRepositoryGateway, never()).add(any(DaHuaVehicleAccessRecordE.class));
assertThrows(IllegalArgumentException.class, () -> applicationService.sync(
"2026-08-01 00:00:00", "2026-08-02 00:00:00", 1001));
}
@Test
@SuppressWarnings("unchecked")
void shouldAdvanceCheckpointOnlyAfterSuccessfulIncrementalSync() {
LocalDateTime checkpoint = LocalDateTime.of(2026, 8, 1, 12, 0, 0);
when(recordRepositoryGateway.findSyncCheckpoint()).thenReturn(checkpoint);
when(daHuaGateway.queryVehicleAccessRecord(anyMap()))
.thenReturn(response(Collections.emptyList()));
applicationService.syncIncremental(100, 5, 10, 1440);
ArgumentCaptor<Map<String, Object>> paramsCaptor = ArgumentCaptor.forClass(Map.class);
verify(daHuaGateway, times(2)).queryVehicleAccessRecord(paramsCaptor.capture());
assertEquals("2026-08-01 11:50:00",
paramsCaptor.getAllValues().get(0).get("enterTimeStrLeft"));
ArgumentCaptor<LocalDateTime> checkpointCaptor = ArgumentCaptor.forClass(LocalDateTime.class);
verify(recordRepositoryGateway).saveSyncCheckpoint(checkpointCaptor.capture());
assertEquals(checkpointCaptor.getValue().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")),
paramsCaptor.getAllValues().get(0).get("enterTimeStrRight"));
}
@Test
void shouldNotAdvanceCheckpointWhenIncrementalSyncContainsInvalidRecord() {
when(daHuaGateway.queryVehicleAccessRecord(anyMap())).thenAnswer(invocation -> {
Map<String, Object> params = invocation.getArgument(0);
return params.containsKey("enterTimeStrLeft")
&& Integer.valueOf(1).equals(params.get("pageNum"))
? response(Collections.singletonList(record(null)))
: response(Collections.emptyList());
});
assertThrows(IllegalStateException.class,
() -> applicationService.syncIncremental(100, 5, 10, 1440));
verify(recordRepositoryGateway, never()).saveSyncCheckpoint(any(LocalDateTime.class));
}
@Test
void shouldFailWhenDahuaRepeatsTheSamePage() {
when(daHuaGateway.queryVehicleAccessRecord(anyMap()))
.thenReturn(response(Collections.singletonList(record("A-1"))));
IllegalStateException exception = assertThrows(IllegalStateException.class,
() -> applicationService.sync(
"2026-08-01 00:00:00", "2026-08-02 00:00:00", 100));
assertTrue(exception.getMessage().contains("repeated page"));
verify(recordRepositoryGateway, times(1)).add(any(DaHuaVehicleAccessRecordE.class));
}
private Map<String, Object> params(String field, int pageNum) {
Map<String, Object> params = new HashMap<>();
params.put(field, "2026-08-01 00:00:00");
params.put("pageNum", pageNum);
params.put("pageSize", 1);
return org.mockito.ArgumentMatchers.argThat(value ->
value.equals(params) || (value.containsKey(field)
&& Integer.valueOf(pageNum).equals(value.get("pageNum"))));
}
private CorpInfoSnapshotE corp(Long id, String name, Integer portArea) {
CorpInfoSnapshotE corp = new CorpInfoSnapshotE();
corp.setId(id);
corp.setCorpName(name);
corp.setPortArea(portArea);
return corp;
}
private Map<String, Object> record(String id) {
Map<String, Object> value = new HashMap<>();
if (id != null) {
value.put("accessId", id);
}
value.put("carNum", "冀A12345");
value.put("enterTimeStr", "2026-08-01 10:00:00");
return value;
}
private Map<String, Object> recordWithDepartmentAndChannel(String id) {
Map<String, Object> value = record(id);
value.put("departmentId", 10L);
value.put("departmentName", "部门");
value.put("enterSluiceDevChnid", "CH-ENTER");
return value;
}
private Map<String, Object> recordWithPlate(String id, String plate) {
Map<String, Object> value = record(id);
value.put("carNum", plate);
return value;
}
private Map<String, Object> response(Collection<Map<String, Object>> records) {
Map<String, Object> data = new HashMap<>();
data.put("pageData", new ArrayList<>(records));
Map<String, Object> response = new HashMap<>();
response.put("passed", true);
response.put("success", true);
response.put("code", "0");
response.put("data", data);
return response;
}
}

View File

@ -31,9 +31,11 @@ public interface DaHuaGateway {
Map<String, Object> queryAccessRecord(DaHuaAccessRecordCmd cmd);
Map<String, Object> queryVehicleAccessRecord(Map<String, Object> params);
String uploadPersonImg(String imageBase64);
Map<String, Object> saveTempVehicle(DaHuaTempVehicleSaveCmd cmd);
Long findDahuaDepartmentIdByCorpId(Long corpId);
}
}

View File

@ -4,7 +4,9 @@ import com.zcloud.primeport.domain.model.CorpInfoSnapshotE;
import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE;
import com.zcloud.primeport.domain.model.DaHuaDeviceE;
import com.zcloud.primeport.domain.model.DaHuaDeviceChannelE;
import com.zcloud.primeport.domain.model.VehicleApplySnapshotE;
import java.util.Collection;
import java.util.List;
public interface DaHuaResourceRepositoryGateway {
@ -25,6 +27,8 @@ public interface DaHuaResourceRepositoryGateway {
List<DaHuaDeviceChannelE> listDeviceChannels();
List<VehicleApplySnapshotE> listVehicleApplySnapshots(Collection<String> licenceNos);
void addDeviceChannel(DaHuaDeviceChannelE channel);
void updateDeviceChannel(DaHuaDeviceChannelE channel);

View File

@ -0,0 +1,20 @@
package com.zcloud.primeport.domain.gateway;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE;
import java.time.LocalDateTime;
import java.util.Collection;
import java.util.List;
public interface DaHuaVehicleAccessRecordRepositoryGateway {
List<DaHuaVehicleAccessRecordE> listByDahuaAccessIds(Collection<String> accessIds);
LocalDateTime findSyncCheckpoint();
void saveSyncCheckpoint(LocalDateTime checkpointTime);
void add(DaHuaVehicleAccessRecordE record);
void update(DaHuaVehicleAccessRecordE record);
}

View File

@ -0,0 +1,11 @@
package com.zcloud.primeport.domain.gateway;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE;
public interface DaHuaVehicleAccessRecordSyncGateway {
DaHuaVehicleAccessRecordSyncResultE syncIncremental(int pageSize, int delayMinutes,
int overlapMinutes, int initialLookbackMinutes);
DaHuaVehicleAccessRecordSyncResultE sync(String startTime, String endTime, int pageSize);
}

View File

@ -0,0 +1,45 @@
package com.zcloud.primeport.domain.model;
import lombok.Data;
import java.time.LocalDateTime;
@Data
public class DaHuaVehicleAccessRecordE {
private Long id;
private String dahuaAccessId;
private Long corpId;
private String corpName;
private Integer portArea;
private Long dahuaDepartmentId;
private String dahuaDepartmentName;
private String dahuaOrgCode;
private String parkingLotCode;
private String parkingLotName;
private Integer carType;
private String carTypeName;
private String carNum;
private String exitCarNum;
private String enterItcChannelCode;
private String enterSluiceChannelCode;
private String enterSluiceChannelName;
private LocalDateTime enterTime;
private String enterImageUrl;
private String exitItcChannelCode;
private String exitSluiceChannelCode;
private String exitSluiceChannelName;
private LocalDateTime exitTime;
private String exitImageUrl;
private String sourceType;
private String rawData;
private LocalDateTime lastSyncTime;
private String deleteEnum;
}

View File

@ -0,0 +1,17 @@
package com.zcloud.primeport.domain.model;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
@Data
public class DaHuaVehicleAccessRecordSyncResultE {
private int pageCount;
private int recordTotal;
private int recordCreated;
private int recordUpdated;
private int recordInvalid;
private List<String> exceptionDetails = new ArrayList<>();
}

View File

@ -0,0 +1,16 @@
package com.zcloud.primeport.domain.model;
import lombok.Data;
/**
*
*/
@Data
public class VehicleApplySnapshotE {
private String licenceNo;
private Long vehicleCorpId;
private String vehicleCorpName;
private Long vehicleDepartmentId;
private String vehicleDepartmentName;
}

View File

@ -0,0 +1,33 @@
package com.zcloud.primeport.dahua.api;
import com.dahuatech.icc.exception.ClientException;
import com.dahuatech.icc.oauth.model.v202010.GeneralResponse;
import com.zcloud.primeport.dahua.client.DaHuaApiClient;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.HashMap;
import java.util.Map;
@Component
public class DaHuaVehicleAccessRecordApi {
private static final String PATH = "/evo-apigw/ipms/caraccess/find/his";
@Autowired(required = false)
private DaHuaApiClient daHuaApiClient;
public boolean isConfigured() {
return daHuaApiClient != null;
}
public GeneralResponse query(Map<String, Object> params) throws ClientException {
if (daHuaApiClient == null) {
return null;
}
String token = daHuaApiClient.getToken();
Map<String, String> headers = new HashMap<>();
headers.put("Authorization", "bearer " + token);
return daHuaApiClient.getForm(PATH, params, headers);
}
}

View File

@ -9,7 +9,7 @@ import org.springframework.context.annotation.Configuration;
@Configuration
@EnableConfigurationProperties({DaHuaProperties.class, DaHuaResourceProperties.class,
DaHuaAccessRecordProperties.class})
DaHuaAccessRecordProperties.class, DaHuaVehicleAccessRecordProperties.class})
public class DaHuaAutoConfig {
@Bean("dahuaApiOauthConfig")

View File

@ -0,0 +1,15 @@
package com.zcloud.primeport.dahua.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
@Data
@ConfigurationProperties(prefix = "dahua.vehicle-access-record")
public class DaHuaVehicleAccessRecordProperties {
private boolean syncEnabled = true;
private Integer syncPageSize = 1000;
private Integer syncDelayMinutes = 5;
private Integer syncOverlapMinutes = 10;
private Integer initialLookbackMinutes = 1440;
}

View File

@ -6,6 +6,7 @@ import cn.hutool.json.JSONUtil;
import com.dahuatech.icc.exception.ClientException;
import com.dahuatech.icc.oauth.model.v202010.GeneralResponse;
import com.zcloud.primeport.dahua.api.DaHuaPersonApi;
import com.zcloud.primeport.dahua.api.DaHuaVehicleAccessRecordApi;
import com.zcloud.primeport.domain.gateway.DaHuaGateway;
import com.zcloud.gbscommon.dahua.cmd.DaHuaAccessRecordCmd;
import com.zcloud.gbscommon.dahua.cmd.DaHuaCarAddCmd;
@ -33,6 +34,9 @@ public class DaHuaGatewayImpl implements DaHuaGateway {
@Autowired(required = false)
private DaHuaPersonApi daHuaPersonApi;
@Autowired(required = false)
private DaHuaVehicleAccessRecordApi daHuaVehicleAccessRecordApi;
@Autowired
private DaHuaDepartmentCorpMappingRepository daHuaDeptMappingRepository;
@ -854,6 +858,39 @@ public class DaHuaGatewayImpl implements DaHuaGateway {
}
}
@Override
public Map<String, Object> queryVehicleAccessRecord(Map<String, Object> params) {
if (daHuaVehicleAccessRecordApi == null || !daHuaVehicleAccessRecordApi.isConfigured()) {
throw new RuntimeException("Dahua service is not configured");
}
Map<String, Object> result = new HashMap<>();
try {
GeneralResponse response = daHuaVehicleAccessRecordApi.query(params);
if (response == null) {
result.put("passed", false);
result.put("success", false);
result.put("code", "SYSTEM_ERROR");
result.put("errMsg", "response is null");
return result;
}
JSONObject json = JSONUtil.parseObj(JSONUtil.toJsonStr(response));
result.put("success", json.getBool("success", false));
result.put("code", json.getStr("code", response.getCode()));
result.put("errMsg", json.getStr("errMsg", ""));
result.put("data", json.get("data"));
result.put("passed", Boolean.TRUE.equals(result.get("success"))
|| "0".equals(String.valueOf(result.get("code")))
|| "200".equals(String.valueOf(result.get("code"))));
return result;
} catch (ClientException e) {
result.put("passed", false);
result.put("success", false);
result.put("code", "CLIENT_EXCEPTION");
result.put("errMsg", e.getMessage());
return result;
}
}
@Override
public String uploadPersonImg(String imageBase64) {
if (daHuaPersonApi == null || !daHuaPersonApi.isConfigured()) {
@ -1202,4 +1239,4 @@ public class DaHuaGatewayImpl implements DaHuaGateway {
}
return mapping.getDahuaDepartmentId();
}
}
}

View File

@ -8,20 +8,25 @@ import com.zcloud.primeport.domain.model.CorpInfoSnapshotE;
import com.zcloud.primeport.domain.model.DaHuaDepartmentCorpMappingE;
import com.zcloud.primeport.domain.model.DaHuaDeviceE;
import com.zcloud.primeport.domain.model.DaHuaDeviceChannelE;
import com.zcloud.primeport.domain.model.VehicleApplySnapshotE;
import com.zcloud.primeport.persistence.dataobject.CorpInfoSnapshotDO;
import com.zcloud.primeport.persistence.dataobject.DaHuaDepartmentCorpMappingDO;
import com.zcloud.primeport.persistence.dataobject.DaHuaDeviceDO;
import com.zcloud.primeport.persistence.dataobject.DaHuaDeviceChannelDO;
import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO;
import com.zcloud.primeport.persistence.repository.CorpInfoSnapshotRepository;
import com.zcloud.primeport.persistence.repository.DaHuaDepartmentCorpMappingRepository;
import com.zcloud.primeport.persistence.repository.DaHuaDeviceRepository;
import com.zcloud.primeport.persistence.repository.DaHuaDeviceChannelRepository;
import com.zcloud.primeport.persistence.repository.VehicleApplyRepository;
import lombok.RequiredArgsConstructor;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
@Service
@ -32,6 +37,7 @@ public class DaHuaResourceRepositoryGatewayImpl implements DaHuaResourceReposito
private final DaHuaDepartmentCorpMappingRepository mappingRepository;
private final DaHuaDeviceRepository deviceRepository;
private final DaHuaDeviceChannelRepository deviceChannelRepository;
private final VehicleApplyRepository vehicleApplyRepository;
private final DefaultIdentifierGenerator idGenerator = new DefaultIdentifierGenerator();
@Override
@ -89,6 +95,15 @@ public class DaHuaResourceRepositoryGatewayImpl implements DaHuaResourceReposito
return convertList(deviceChannelRepository.list(new LambdaQueryWrapper<>()), DaHuaDeviceChannelE.class);
}
@Override
public List<VehicleApplySnapshotE> listVehicleApplySnapshots(Collection<String> licenceNos) {
if (licenceNos == null || licenceNos.isEmpty()) {
return Collections.emptyList();
}
List<VehicleApplySnapshotDO> snapshots = vehicleApplyRepository.listAccessMatchSnapshots(licenceNos);
return convertList(snapshots, VehicleApplySnapshotE.class);
}
@Override
public void addDeviceChannel(DaHuaDeviceChannelE channel) {
DaHuaDeviceChannelDO data = copy(channel, DaHuaDeviceChannelDO.class);

View File

@ -0,0 +1,90 @@
package com.zcloud.primeport.gatewayimpl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.incrementer.DefaultIdentifierGenerator;
import com.jjb.saas.framework.repository.basedo.BaseDO;
import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordRepositoryGateway;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordE;
import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO;
import com.zcloud.primeport.persistence.repository.DaHuaVehicleAccessRecordRepository;
import lombok.RequiredArgsConstructor;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
@Service
@RequiredArgsConstructor
public class DaHuaVehicleAccessRecordRepositoryGatewayImpl
implements DaHuaVehicleAccessRecordRepositoryGateway {
private final DaHuaVehicleAccessRecordRepository repository;
private final DefaultIdentifierGenerator idGenerator = new DefaultIdentifierGenerator();
@Override
public List<DaHuaVehicleAccessRecordE> listByDahuaAccessIds(Collection<String> accessIds) {
if (accessIds == null || accessIds.isEmpty()) {
return Collections.emptyList();
}
LambdaQueryWrapper<DaHuaVehicleAccessRecordDO> query = new LambdaQueryWrapper<>();
query.in(DaHuaVehicleAccessRecordDO::getDahuaAccessId, accessIds);
return convertList(repository.list(query));
}
@Override
public LocalDateTime findSyncCheckpoint() {
return repository.findSyncCheckpoint();
}
@Override
public void saveSyncCheckpoint(LocalDateTime checkpointTime) {
repository.saveSyncCheckpoint(checkpointTime);
}
@Override
public void add(DaHuaVehicleAccessRecordE record) {
DaHuaVehicleAccessRecordDO data = copy(record, DaHuaVehicleAccessRecordDO.class);
initializeBase(data);
repository.upsert(data);
record.setId(data.getId());
}
@Override
public void update(DaHuaVehicleAccessRecordE record) {
DaHuaVehicleAccessRecordDO data = copy(record, DaHuaVehicleAccessRecordDO.class);
data.setUpdateTime(LocalDateTime.now());
repository.upsert(data);
}
private void initializeBase(BaseDO target) {
LocalDateTime now = LocalDateTime.now();
target.setId(idGenerator.nextId(target).longValue());
target.setDeleteEnum("FALSE");
target.setEnv("PROD");
target.setVersion(0);
target.setCreateTime(now);
target.setUpdateTime(now);
}
private List<DaHuaVehicleAccessRecordE> convertList(List<DaHuaVehicleAccessRecordDO> source) {
List<DaHuaVehicleAccessRecordE> result = new ArrayList<>();
for (DaHuaVehicleAccessRecordDO item : source) {
result.add(copy(item, DaHuaVehicleAccessRecordE.class));
}
return result;
}
private <T> T copy(Object source, Class<T> targetType) {
try {
T target = targetType.newInstance();
BeanUtils.copyProperties(source, target);
return target;
} catch (InstantiationException | IllegalAccessException e) {
throw new IllegalStateException("Failed to convert vehicle access record", e);
}
}
}

View File

@ -0,0 +1,41 @@
package com.zcloud.primeport.persistence.dataobject;
import com.baomidou.mybatisplus.annotation.TableName;
import com.jjb.saas.framework.repository.basedo.BaseDO;
import lombok.Data;
import lombok.EqualsAndHashCode;
import java.time.LocalDateTime;
@Data
@TableName("dahua_vehicle_access_record")
@EqualsAndHashCode(callSuper = true)
public class DaHuaVehicleAccessRecordDO extends BaseDO {
private String dahuaAccessId;
private Long corpId;
private String corpName;
private Integer portArea;
private Long dahuaDepartmentId;
private String dahuaDepartmentName;
private String dahuaOrgCode;
private String parkingLotCode;
private String parkingLotName;
private Integer carType;
private String carTypeName;
private String carNum;
private String exitCarNum;
private String enterItcChannelCode;
private String enterSluiceChannelCode;
private String enterSluiceChannelName;
private LocalDateTime enterTime;
private String enterImageUrl;
private String exitItcChannelCode;
private String exitSluiceChannelCode;
private String exitSluiceChannelName;
private LocalDateTime exitTime;
private String exitImageUrl;
private String sourceType;
private String rawData;
private LocalDateTime lastSyncTime;
}

View File

@ -0,0 +1,15 @@
package com.zcloud.primeport.persistence.dataobject;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
@Data
@TableName("vehicle_apply")
public class VehicleApplySnapshotDO {
private String licenceNo;
private Long vehicleCorpId;
private String vehicleCorpName;
private Long vehicleDepartmentId;
private String vehicleDepartmentName;
}

View File

@ -0,0 +1,18 @@
package com.zcloud.primeport.persistence.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.time.LocalDateTime;
@Mapper
public interface DaHuaVehicleAccessRecordMapper extends BaseMapper<DaHuaVehicleAccessRecordDO> {
LocalDateTime findSyncCheckpoint();
int upsertSyncCheckpoint(@Param("checkpointTime") LocalDateTime checkpointTime);
int upsert(@Param("record") DaHuaVehicleAccessRecordDO record);
}

View File

@ -9,6 +9,7 @@ import com.jjb.saas.framework.datascope.annotation.DataScopes;
import com.zcloud.primeport.persistence.dataobject.AppCountDTO;
import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto;
import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO;
import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO;
import com.zcloud.primeport.plan.mjDevice.dto.OneLevelCarSaveDto;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
@ -16,6 +17,7 @@ import org.apache.ibatis.annotations.Param;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Collection;
/**
* web-infrastructure
@ -40,5 +42,7 @@ public interface VehicleApplyMapper extends BaseMapper<VehicleApplyDO> {
List<VehicleApplyDO> getAppCount(@Param("params") Map<String, Object> parmas);
List<OneLevelCarSaveDto> listTodayHasPort(String licenceNo);
List<VehicleApplySnapshotDO> listAccessMatchSnapshots(@Param("licenceNos") Collection<String> licenceNos);
}

View File

@ -0,0 +1,15 @@
package com.zcloud.primeport.persistence.repository;
import com.jjb.saas.framework.repository.repo.BaseRepository;
import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO;
import java.time.LocalDateTime;
public interface DaHuaVehicleAccessRecordRepository extends BaseRepository<DaHuaVehicleAccessRecordDO> {
LocalDateTime findSyncCheckpoint();
void saveSyncCheckpoint(LocalDateTime checkpointTime);
void upsert(DaHuaVehicleAccessRecordDO record);
}

View File

@ -5,11 +5,13 @@ import com.jjb.saas.framework.repository.repo.BaseRepository;
import com.zcloud.primeport.persistence.dataobject.AppCountDTO;
import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto;
import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO;
import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO;
import com.zcloud.primeport.plan.mjDevice.dto.OneLevelCarSaveDto;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Collection;
/**
* web-infrastructure
@ -32,5 +34,7 @@ public interface VehicleApplyRepository extends BaseRepository<VehicleApplyDO> {
List<AppCountDTO> getAppCount(Map<String, Object> parmas);
List<OneLevelCarSaveDto> listTodayHasPort(String licenceNo);
List<VehicleApplySnapshotDO> listAccessMatchSnapshots(Collection<String> licenceNos);
}

View File

@ -0,0 +1,30 @@
package com.zcloud.primeport.persistence.repository.impl;
import com.jjb.saas.framework.repository.repo.impl.BaseRepositoryImpl;
import com.zcloud.primeport.persistence.dataobject.DaHuaVehicleAccessRecordDO;
import com.zcloud.primeport.persistence.mapper.DaHuaVehicleAccessRecordMapper;
import com.zcloud.primeport.persistence.repository.DaHuaVehicleAccessRecordRepository;
import org.springframework.stereotype.Repository;
import java.time.LocalDateTime;
@Repository
public class DaHuaVehicleAccessRecordRepositoryImpl
extends BaseRepositoryImpl<DaHuaVehicleAccessRecordMapper, DaHuaVehicleAccessRecordDO>
implements DaHuaVehicleAccessRecordRepository {
@Override
public LocalDateTime findSyncCheckpoint() {
return baseMapper.findSyncCheckpoint();
}
@Override
public void saveSyncCheckpoint(LocalDateTime checkpointTime) {
baseMapper.upsertSyncCheckpoint(checkpointTime);
}
@Override
public void upsert(DaHuaVehicleAccessRecordDO record) {
baseMapper.upsert(record);
}
}

View File

@ -16,6 +16,7 @@ import com.zcloud.primeport.domain.enums.VehicleBelongTypeEnum;
import com.zcloud.primeport.persistence.dataobject.AppCountDTO;
import com.zcloud.primeport.persistence.dataobject.FgsVehicleCountDto;
import com.zcloud.primeport.persistence.dataobject.VehicleApplyDO;
import com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO;
import com.zcloud.primeport.persistence.dataobject.VehicleMessageDO;
import com.zcloud.primeport.persistence.mapper.VehicleApplyMapper;
import com.zcloud.primeport.persistence.repository.VehicleApplyRepository;
@ -112,5 +113,13 @@ public class VehicleApplyRepositoryImpl extends BaseRepositoryImpl<VehicleApplyM
public List<OneLevelCarSaveDto> listTodayHasPort(String licenceNo) {
return vehicleApplyMapper.listTodayHasPort(licenceNo);
}
@Override
public List<VehicleApplySnapshotDO> listAccessMatchSnapshots(Collection<String> licenceNos) {
if (licenceNos == null || licenceNos.isEmpty()) {
return Collections.emptyList();
}
return vehicleApplyMapper.listAccessMatchSnapshots(licenceNos);
}
}

View File

@ -0,0 +1,43 @@
package com.zcloud.primeport.plan;
import com.jjb.saas.framework.job.Job;
import com.jjb.saas.framework.job.annotation.JobRegister;
import com.xxl.job.core.biz.model.ReturnT;
import com.xxl.job.core.handler.annotation.XxlJob;
import com.zcloud.primeport.dahua.config.DaHuaVehicleAccessRecordProperties;
import com.zcloud.primeport.domain.gateway.DaHuaVehicleAccessRecordSyncGateway;
import com.zcloud.primeport.domain.model.DaHuaVehicleAccessRecordSyncResultE;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
@Component
@RequiredArgsConstructor
@Slf4j
public class DaHuaVehicleAccessRecordSyncXxlJob implements Job {
private final DaHuaVehicleAccessRecordSyncGateway syncGateway;
private final DaHuaVehicleAccessRecordProperties properties;
@Override
@JobRegister(cron = "0 */5 * * * ?", jobDesc = "大华车辆进出记录增量同步", triggerStatus = 1)
@XxlJob("com.zcloud.plan.DaHuaVehicleAccessRecordSyncXxlJob")
public ReturnT<String> execute(String param) {
if (!properties.isSyncEnabled()) {
log.info("大华车辆进出记录增量同步已关闭");
return ReturnT.SUCCESS;
}
try {
DaHuaVehicleAccessRecordSyncResultE result = syncGateway.syncIncremental(
properties.getSyncPageSize(), properties.getSyncDelayMinutes(),
properties.getSyncOverlapMinutes(), properties.getInitialLookbackMinutes());
log.info("大华车辆进出记录增量同步完成: pages={}, total={}, created={}, updated={}, invalid={}",
result.getPageCount(), result.getRecordTotal(), result.getRecordCreated(),
result.getRecordUpdated(), result.getRecordInvalid());
return ReturnT.SUCCESS;
} catch (Exception e) {
log.error("大华车辆进出记录增量同步失败", e);
return ReturnT.FAIL;
}
}
}

View File

@ -498,3 +498,62 @@ CREATE TABLE `dahua_access_record` (
KEY `idx_dahua_access_record_person_time` (`person_code`, `swing_time`),
KEY `idx_dahua_access_record_create_time` (`dahua_create_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华门禁通行记录表';
CREATE TABLE `dahua_vehicle_access_record` (
`id` bigint NOT NULL COMMENT '主键ID',
`dahua_access_id` varchar(128) NOT NULL COMMENT '大华车辆进出记录accessId幂等键',
`corp_id` bigint DEFAULT NULL COMMENT '平台企业ID快照',
`corp_name` varchar(255) DEFAULT NULL COMMENT '平台企业名称快照',
`port_area` int DEFAULT NULL COMMENT '企业港区快照',
`dahua_department_id` bigint DEFAULT NULL COMMENT '大华部门ID',
`dahua_department_name` varchar(255) DEFAULT NULL COMMENT '大华部门名称',
`dahua_org_code` varchar(64) DEFAULT NULL COMMENT '大华组织编码',
`parking_lot_code` varchar(128) DEFAULT NULL COMMENT '停车场编码',
`parking_lot_name` varchar(255) DEFAULT NULL COMMENT '停车场名称',
`car_type` int DEFAULT NULL COMMENT '车辆类型编码',
`car_type_name` varchar(128) DEFAULT NULL COMMENT '车辆类型名称',
`car_num` varchar(64) DEFAULT NULL COMMENT '入场车牌',
`exit_car_num` varchar(64) DEFAULT NULL COMMENT '出场车牌',
`enter_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '入场卡口相机通道编码',
`enter_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '入场道闸通道编码',
`enter_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '入场道闸通道名称',
`enter_time` datetime DEFAULT NULL COMMENT '入场时间',
`enter_image_url` varchar(1024) DEFAULT NULL COMMENT '入场图片',
`exit_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '出场卡口相机通道编码',
`exit_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '出场道闸通道编码',
`exit_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '出场道闸通道名称',
`exit_time` datetime DEFAULT NULL COMMENT '出场时间',
`exit_image_url` varchar(1024) DEFAULT NULL COMMENT '出场图片',
`source_type` varchar(32) NOT NULL DEFAULT 'PULL' COMMENT '数据来源',
`raw_data` json DEFAULT NULL COMMENT '大华原始记录JSON',
`last_sync_time` datetime DEFAULT NULL COMMENT '最近同步时间',
`delete_enum` varchar(32) NOT NULL DEFAULT 'FALSE' COMMENT '删除标识',
`create_time` datetime DEFAULT NULL,
`update_time` datetime DEFAULT NULL,
`create_id` bigint DEFAULT NULL,
`update_id` bigint DEFAULT NULL,
`env` varchar(32) NOT NULL DEFAULT 'PROD',
`create_name` varchar(255) DEFAULT NULL,
`update_name` varchar(255) DEFAULT NULL,
`tenant_id` bigint DEFAULT NULL,
`org_id` bigint DEFAULT NULL,
`version` int NOT NULL DEFAULT '0',
`remarks` varchar(500) DEFAULT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_dahua_vehicle_access_record_access_id` (`dahua_access_id`),
KEY `idx_dahua_vehicle_access_record_corp_enter` (`corp_id`, `enter_time`),
KEY `idx_dahua_vehicle_access_record_corp_exit` (`corp_id`, `exit_time`),
KEY `idx_dahua_vehicle_access_record_car_enter` (`car_num`, `enter_time`),
KEY `idx_dahua_vehicle_access_record_car_exit` (`exit_car_num`, `exit_time`),
KEY `idx_dahua_vehicle_access_record_department` (`dahua_department_id`),
KEY `idx_dahua_vehicle_access_record_enter_channel` (`enter_sluice_channel_code`),
KEY `idx_dahua_vehicle_access_record_exit_channel` (`exit_sluice_channel_code`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华车辆进出记录表';
CREATE TABLE `dahua_sync_checkpoint` (
`sync_type` varchar(64) NOT NULL COMMENT '同步任务类型',
`checkpoint_time` datetime NOT NULL COMMENT '最近完整成功同步窗口的结束时间',
`create_time` datetime DEFAULT NULL COMMENT '创建时间',
`update_time` datetime DEFAULT NULL COMMENT '更新时间',
PRIMARY KEY (`sync_type`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华同步检查点';

View File

@ -0,0 +1,58 @@
CREATE TABLE IF NOT EXISTS `dahua_vehicle_access_record` (
`id` bigint NOT NULL COMMENT '主键ID',
`dahua_access_id` varchar(128) NOT NULL COMMENT '大华车辆进出记录accessId幂等键',
`corp_id` bigint DEFAULT NULL COMMENT '平台企业ID快照',
`corp_name` varchar(255) DEFAULT NULL COMMENT '平台企业名称快照',
`port_area` int DEFAULT NULL COMMENT '企业港区快照',
`dahua_department_id` bigint DEFAULT NULL COMMENT '大华部门ID',
`dahua_department_name` varchar(255) DEFAULT NULL COMMENT '大华部门名称',
`dahua_org_code` varchar(64) DEFAULT NULL COMMENT '大华组织编码',
`parking_lot_code` varchar(128) DEFAULT NULL COMMENT '停车场编码',
`parking_lot_name` varchar(255) DEFAULT NULL COMMENT '停车场名称',
`car_type` int DEFAULT NULL COMMENT '车辆类型编码',
`car_type_name` varchar(128) DEFAULT NULL COMMENT '车辆类型名称',
`car_num` varchar(64) DEFAULT NULL COMMENT '入场车牌',
`exit_car_num` varchar(64) DEFAULT NULL COMMENT '出场车牌',
`enter_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '入场卡口相机通道编码',
`enter_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '入场道闸通道编码',
`enter_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '入场道闸通道名称',
`enter_time` datetime DEFAULT NULL COMMENT '入场时间',
`enter_image_url` varchar(1024) DEFAULT NULL COMMENT '入场图片',
`exit_itc_channel_code` varchar(128) DEFAULT NULL COMMENT '出场卡口相机通道编码',
`exit_sluice_channel_code` varchar(128) DEFAULT NULL COMMENT '出场道闸通道编码',
`exit_sluice_channel_name` varchar(255) DEFAULT NULL COMMENT '出场道闸通道名称',
`exit_time` datetime DEFAULT NULL COMMENT '出场时间',
`exit_image_url` varchar(1024) DEFAULT NULL COMMENT '出场图片',
`source_type` varchar(32) NOT NULL DEFAULT 'PULL' COMMENT '数据来源',
`raw_data` json DEFAULT NULL COMMENT '大华原始记录JSON',
`last_sync_time` datetime DEFAULT NULL COMMENT '最近同步时间',
`delete_enum` varchar(32) NOT NULL DEFAULT 'FALSE' COMMENT '删除标识',
`create_time` datetime DEFAULT NULL,
`update_time` datetime DEFAULT NULL,
`create_id` bigint DEFAULT NULL,
`update_id` bigint DEFAULT NULL,
`env` varchar(32) NOT NULL DEFAULT 'PROD',
`create_name` varchar(255) DEFAULT NULL,
`update_name` varchar(255) DEFAULT NULL,
`tenant_id` bigint DEFAULT NULL,
`org_id` bigint DEFAULT NULL,
`version` int NOT NULL DEFAULT '0',
`remarks` varchar(500) DEFAULT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_dahua_vehicle_access_record_access_id` (`dahua_access_id`),
KEY `idx_dahua_vehicle_access_record_corp_enter` (`corp_id`, `enter_time`),
KEY `idx_dahua_vehicle_access_record_corp_exit` (`corp_id`, `exit_time`),
KEY `idx_dahua_vehicle_access_record_car_enter` (`car_num`, `enter_time`),
KEY `idx_dahua_vehicle_access_record_car_exit` (`exit_car_num`, `exit_time`),
KEY `idx_dahua_vehicle_access_record_department` (`dahua_department_id`),
KEY `idx_dahua_vehicle_access_record_enter_channel` (`enter_sluice_channel_code`),
KEY `idx_dahua_vehicle_access_record_exit_channel` (`exit_sluice_channel_code`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华车辆进出记录表';
CREATE TABLE IF NOT EXISTS `dahua_sync_checkpoint` (
`sync_type` varchar(64) NOT NULL COMMENT '同步任务类型',
`checkpoint_time` datetime NOT NULL COMMENT '最近完整成功同步窗口的结束时间',
`create_time` datetime DEFAULT NULL COMMENT '创建时间',
`update_time` datetime DEFAULT NULL COMMENT '更新时间',
PRIMARY KEY (`sync_type`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='大华同步检查点';

View File

@ -0,0 +1,75 @@
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
"http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.zcloud.primeport.persistence.mapper.DaHuaVehicleAccessRecordMapper">
<select id="findSyncCheckpoint" resultType="java.time.LocalDateTime">
SELECT checkpoint_time
FROM dahua_sync_checkpoint
WHERE sync_type = 'VEHICLE_ACCESS_RECORD'
</select>
<insert id="upsertSyncCheckpoint">
INSERT INTO dahua_sync_checkpoint (sync_type, checkpoint_time, create_time, update_time)
VALUES ('VEHICLE_ACCESS_RECORD', #{checkpointTime}, NOW(), NOW())
ON DUPLICATE KEY UPDATE
checkpoint_time = GREATEST(checkpoint_time, VALUES(checkpoint_time)),
update_time = NOW()
</insert>
<insert id="upsert">
INSERT INTO dahua_vehicle_access_record (
id, dahua_access_id, corp_id, corp_name, port_area,
dahua_department_id, dahua_department_name, dahua_org_code,
parking_lot_code, parking_lot_name, car_type, car_type_name,
car_num, exit_car_num,
enter_itc_channel_code, enter_sluice_channel_code, enter_sluice_channel_name,
enter_time, enter_image_url,
exit_itc_channel_code, exit_sluice_channel_code, exit_sluice_channel_name,
exit_time, exit_image_url,
source_type, raw_data, last_sync_time,
delete_enum, create_time, update_time, env, version
) VALUES (
#{record.id}, #{record.dahuaAccessId}, #{record.corpId}, #{record.corpName}, #{record.portArea},
#{record.dahuaDepartmentId}, #{record.dahuaDepartmentName}, #{record.dahuaOrgCode},
#{record.parkingLotCode}, #{record.parkingLotName}, #{record.carType}, #{record.carTypeName},
#{record.carNum}, #{record.exitCarNum},
#{record.enterItcChannelCode}, #{record.enterSluiceChannelCode}, #{record.enterSluiceChannelName},
#{record.enterTime}, #{record.enterImageUrl},
#{record.exitItcChannelCode}, #{record.exitSluiceChannelCode}, #{record.exitSluiceChannelName},
#{record.exitTime}, #{record.exitImageUrl},
#{record.sourceType}, #{record.rawData}, #{record.lastSyncTime},
#{record.deleteEnum}, #{record.createTime}, #{record.updateTime}, #{record.env}, #{record.version}
)
ON DUPLICATE KEY UPDATE
corp_id = COALESCE(VALUES(corp_id), corp_id),
corp_name = COALESCE(VALUES(corp_name), corp_name),
port_area = COALESCE(VALUES(port_area), port_area),
dahua_department_id = COALESCE(VALUES(dahua_department_id), dahua_department_id),
dahua_department_name = COALESCE(VALUES(dahua_department_name), dahua_department_name),
dahua_org_code = COALESCE(VALUES(dahua_org_code), dahua_org_code),
parking_lot_code = COALESCE(VALUES(parking_lot_code), parking_lot_code),
parking_lot_name = COALESCE(VALUES(parking_lot_name), parking_lot_name),
car_type = COALESCE(VALUES(car_type), car_type),
car_type_name = COALESCE(VALUES(car_type_name), car_type_name),
car_num = COALESCE(VALUES(car_num), car_num),
exit_car_num = COALESCE(VALUES(exit_car_num), exit_car_num),
enter_itc_channel_code = COALESCE(VALUES(enter_itc_channel_code), enter_itc_channel_code),
enter_sluice_channel_code = COALESCE(VALUES(enter_sluice_channel_code), enter_sluice_channel_code),
enter_sluice_channel_name = COALESCE(VALUES(enter_sluice_channel_name), enter_sluice_channel_name),
enter_time = COALESCE(VALUES(enter_time), enter_time),
enter_image_url = COALESCE(VALUES(enter_image_url), enter_image_url),
exit_itc_channel_code = COALESCE(VALUES(exit_itc_channel_code), exit_itc_channel_code),
exit_sluice_channel_code = COALESCE(VALUES(exit_sluice_channel_code), exit_sluice_channel_code),
exit_sluice_channel_name = COALESCE(VALUES(exit_sluice_channel_name), exit_sluice_channel_name),
exit_time = COALESCE(VALUES(exit_time), exit_time),
exit_image_url = COALESCE(VALUES(exit_image_url), exit_image_url),
source_type = VALUES(source_type),
raw_data = VALUES(raw_data),
last_sync_time = VALUES(last_sync_time),
delete_enum = VALUES(delete_enum),
update_time = VALUES(update_time)
</insert>
</mapper>

View File

@ -220,5 +220,20 @@
licence_no,
licence_type
</select>
<select id="listAccessMatchSnapshots"
resultType="com.zcloud.primeport.persistence.dataobject.VehicleApplySnapshotDO">
SELECT licence_no,
vehicle_corp_id,
vehicle_corp_name,
vehicle_department_id,
vehicle_department_name
FROM vehicle_apply
WHERE delete_enum = 'FALSE'
AND licence_no IN
<foreach collection="licenceNos" item="licenceNo" open="(" separator="," close=")">
#{licenceNo}
</foreach>
ORDER BY id DESC
</select>
</mapper>