feat(community): improve metering data upload with batch processing and deduplication
This commit is contained in:
parent
85d8bc0011
commit
5ff22e53bd
@ -3,6 +3,7 @@ package at.mueller.eeg.backend.community.repository;
|
|||||||
import at.mueller.eeg.backend.community.domain.MeteringData;
|
import at.mueller.eeg.backend.community.domain.MeteringData;
|
||||||
import at.mueller.eeg.backend.community.domain.MeteringDataType;
|
import at.mueller.eeg.backend.community.domain.MeteringDataType;
|
||||||
import org.springframework.data.jpa.repository.JpaRepository;
|
import org.springframework.data.jpa.repository.JpaRepository;
|
||||||
|
import org.springframework.data.jpa.repository.Modifying;
|
||||||
import org.springframework.data.jpa.repository.Query;
|
import org.springframework.data.jpa.repository.Query;
|
||||||
import org.springframework.data.repository.query.Param;
|
import org.springframework.data.repository.query.Param;
|
||||||
|
|
||||||
@ -30,6 +31,8 @@ public interface MeteringDataRepository extends JpaRepository<MeteringData, UUID
|
|||||||
@Param("dataType") MeteringDataType dataType,
|
@Param("dataType") MeteringDataType dataType,
|
||||||
@Param("intervalStarts") List<LocalDateTime> intervalStarts);
|
@Param("intervalStarts") List<LocalDateTime> intervalStarts);
|
||||||
|
|
||||||
|
List<MeteringData> findByMeteringPointId(UUID meteringPointId);
|
||||||
|
|
||||||
@Query("SELECT FUNCTION('DAYOFYEAR', m.intervalStart), FUNCTION('YEAR', m.intervalStart), SUM(m.kwh) " +
|
@Query("SELECT FUNCTION('DAYOFYEAR', m.intervalStart), FUNCTION('YEAR', m.intervalStart), SUM(m.kwh) " +
|
||||||
"FROM MeteringData m WHERE m.meteringPoint.id = :meteringPointId AND m.dataType = 'CONSUMPTION' " +
|
"FROM MeteringData m WHERE m.meteringPoint.id = :meteringPointId AND m.dataType = 'CONSUMPTION' " +
|
||||||
"AND m.intervalStart BETWEEN :from AND :to " +
|
"AND m.intervalStart BETWEEN :from AND :to " +
|
||||||
|
|||||||
@ -9,6 +9,7 @@ import at.mueller.eeg.backend.community.mapper.MeteringDataMapper;
|
|||||||
import at.mueller.eeg.backend.community.repository.MeteringDataRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringDataRepository;
|
||||||
import at.mueller.eeg.backend.community.repository.MeteringDataUploadRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringDataUploadRepository;
|
||||||
import at.mueller.eeg.backend.community.repository.MeteringPointRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringPointRepository;
|
||||||
|
import jakarta.persistence.EntityManager;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import org.springframework.transaction.annotation.Transactional;
|
import org.springframework.transaction.annotation.Transactional;
|
||||||
@ -26,6 +27,7 @@ public class MeteringDataService {
|
|||||||
private final MeteringDataUploadRepository uploadRepository;
|
private final MeteringDataUploadRepository uploadRepository;
|
||||||
private final MeteringPointRepository meteringPointRepository;
|
private final MeteringPointRepository meteringPointRepository;
|
||||||
private final MeteringDataMapper mapper;
|
private final MeteringDataMapper mapper;
|
||||||
|
private final EntityManager entityManager;
|
||||||
|
|
||||||
@Transactional
|
@Transactional
|
||||||
public MeteringDataUploadResponse uploadMeteringData(MeteringDataUploadRequest request, UUID uploadedBy) {
|
public MeteringDataUploadResponse uploadMeteringData(MeteringDataUploadRequest request, UUID uploadedBy) {
|
||||||
@ -43,18 +45,42 @@ public class MeteringDataService {
|
|||||||
upload.setStatus(UploadStatus.PROCESSING);
|
upload.setStatus(UploadStatus.PROCESSING);
|
||||||
upload.setUploadedBy(uploadedBy);
|
upload.setUploadedBy(uploadedBy);
|
||||||
MeteringDataUpload savedUpload = uploadRepository.save(upload);
|
MeteringDataUpload savedUpload = uploadRepository.save(upload);
|
||||||
|
entityManager.flush();
|
||||||
|
|
||||||
List<MeteringData> entities = request.records().stream()
|
// Convert all records to entities first
|
||||||
|
List<MeteringData> allEntities = request.records().stream()
|
||||||
.map(record -> toEntity(record, pointsByAtNumber.get(record.atNumber()),
|
.map(record -> toEntity(record, pointsByAtNumber.get(record.atNumber()),
|
||||||
savedUpload.getId(), request.source()))
|
savedUpload.getId(), request.source()))
|
||||||
.toList();
|
.toList();
|
||||||
|
|
||||||
deduplicateAndSave(entities);
|
// Deduplicate across all records first
|
||||||
|
List<MeteringData> deduplicated = deduplicate(allEntities);
|
||||||
|
|
||||||
|
// Delete existing records for all metering points involved
|
||||||
|
Set<UUID> meteringPointIds = deduplicated.stream()
|
||||||
|
.map(e -> e.getMeteringPoint().getId())
|
||||||
|
.collect(Collectors.toSet());
|
||||||
|
for (UUID pointId : meteringPointIds) {
|
||||||
|
List<MeteringData> existing = meteringDataRepository.findByMeteringPointId(pointId);
|
||||||
|
if (!existing.isEmpty()) {
|
||||||
|
meteringDataRepository.deleteAll(existing);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
entityManager.flush();
|
||||||
|
entityManager.clear();
|
||||||
|
|
||||||
|
// Save in batches of 200
|
||||||
|
int batchSize = 200;
|
||||||
|
for (int i = 0; i < deduplicated.size(); i += batchSize) {
|
||||||
|
List<MeteringData> batch = deduplicated.subList(i, Math.min(i + batchSize, deduplicated.size()));
|
||||||
|
meteringDataRepository.saveAll(batch);
|
||||||
|
}
|
||||||
|
entityManager.flush();
|
||||||
|
|
||||||
savedUpload.setStatus(UploadStatus.COMPLETED);
|
savedUpload.setStatus(UploadStatus.COMPLETED);
|
||||||
uploadRepository.save(savedUpload);
|
uploadRepository.save(savedUpload);
|
||||||
|
|
||||||
return MeteringDataUploadResponse.success(savedUpload.getId(), entities.size());
|
return MeteringDataUploadResponse.success(savedUpload.getId(), deduplicated.size());
|
||||||
}
|
}
|
||||||
|
|
||||||
public List<MeteringDataResponse> getMeteringData(UUID userId, UUID meteringPointId,
|
public List<MeteringDataResponse> getMeteringData(UUID userId, UUID meteringPointId,
|
||||||
@ -152,36 +178,19 @@ public class MeteringDataService {
|
|||||||
return entity;
|
return entity;
|
||||||
}
|
}
|
||||||
|
|
||||||
private void deduplicateAndSave(List<MeteringData> entities) {
|
private List<MeteringData> deduplicate(List<MeteringData> entities) {
|
||||||
if (entities.isEmpty()) {
|
if (entities.isEmpty()) {
|
||||||
return;
|
return entities;
|
||||||
}
|
}
|
||||||
|
|
||||||
Map<UUID, Map<MeteringDataType, List<LocalDateTime>>> grouped = entities.stream()
|
// Deduplicate - keep last record per (pointId, dataType, intervalStart)
|
||||||
.collect(Collectors.groupingBy(
|
Map<String, MeteringData> uniqueByKey = new LinkedHashMap<>();
|
||||||
e -> e.getMeteringPoint().getId(),
|
for (MeteringData entity : entities) {
|
||||||
Collectors.groupingBy(
|
String key = entity.getMeteringPoint().getId() + "|" +
|
||||||
MeteringData::getDataType,
|
entity.getDataType() + "|" +
|
||||||
Collectors.mapping(MeteringData::getIntervalStart, Collectors.toList())
|
entity.getIntervalStart();
|
||||||
)
|
uniqueByKey.put(key, entity);
|
||||||
));
|
|
||||||
|
|
||||||
Set<MeteringData> toDelete = new LinkedHashSet<>();
|
|
||||||
|
|
||||||
for (var pointEntry : grouped.entrySet()) {
|
|
||||||
UUID pointId = pointEntry.getKey();
|
|
||||||
for (var typeEntry : pointEntry.getValue().entrySet()) {
|
|
||||||
MeteringDataType dataType = typeEntry.getKey();
|
|
||||||
List<LocalDateTime> starts = typeEntry.getValue();
|
|
||||||
toDelete.addAll(meteringDataRepository.findByMeteringPointIdAndDataTypeAndIntervalStartIn(
|
|
||||||
pointId, dataType, starts));
|
|
||||||
}
|
}
|
||||||
}
|
return new ArrayList<>(uniqueByKey.values());
|
||||||
|
|
||||||
if (!toDelete.isEmpty()) {
|
|
||||||
meteringDataRepository.deleteAll(new ArrayList<>(toDelete));
|
|
||||||
}
|
|
||||||
|
|
||||||
meteringDataRepository.saveAll(entities);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -9,6 +9,7 @@ import at.mueller.eeg.backend.community.mapper.MeteringDataMapper;
|
|||||||
import at.mueller.eeg.backend.community.repository.MeteringDataRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringDataRepository;
|
||||||
import at.mueller.eeg.backend.community.repository.MeteringDataUploadRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringDataUploadRepository;
|
||||||
import at.mueller.eeg.backend.community.repository.MeteringPointRepository;
|
import at.mueller.eeg.backend.community.repository.MeteringPointRepository;
|
||||||
|
import jakarta.persistence.EntityManager;
|
||||||
import org.junit.jupiter.api.BeforeEach;
|
import org.junit.jupiter.api.BeforeEach;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
import org.junit.jupiter.api.extension.ExtendWith;
|
import org.junit.jupiter.api.extension.ExtendWith;
|
||||||
@ -18,6 +19,8 @@ import org.mockito.junit.jupiter.MockitoExtension;
|
|||||||
import org.springframework.security.access.AccessDeniedException;
|
import org.springframework.security.access.AccessDeniedException;
|
||||||
|
|
||||||
import java.time.LocalDateTime;
|
import java.time.LocalDateTime;
|
||||||
|
import java.time.temporal.ChronoUnit;
|
||||||
|
import java.util.ArrayList;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
@ -42,6 +45,9 @@ class MeteringDataServiceTest {
|
|||||||
@Mock
|
@Mock
|
||||||
private MeteringDataMapper mapper;
|
private MeteringDataMapper mapper;
|
||||||
|
|
||||||
|
@Mock
|
||||||
|
private EntityManager entityManager;
|
||||||
|
|
||||||
@InjectMocks
|
@InjectMocks
|
||||||
private MeteringDataService meteringDataService;
|
private MeteringDataService meteringDataService;
|
||||||
|
|
||||||
@ -390,4 +396,132 @@ class MeteringDataServiceTest {
|
|||||||
assertEquals(UploadStatus.COMPLETED, response.status());
|
assertEquals(UploadStatus.COMPLETED, response.status());
|
||||||
assertEquals(1, response.recordCount());
|
assertEquals(1, response.recordCount());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void uploadMeteringData_largeBatchOver1000Records() {
|
||||||
|
when(meteringPointRepository.findByAtNumberIn(any())).thenReturn(List.of(testMeteringPoint));
|
||||||
|
when(uploadRepository.save(any())).thenAnswer(invocation -> {
|
||||||
|
MeteringDataUpload upload = invocation.getArgument(0);
|
||||||
|
upload.setId(UUID.randomUUID());
|
||||||
|
return upload;
|
||||||
|
});
|
||||||
|
when(meteringDataRepository.saveAll(any())).thenAnswer(invocation -> invocation.getArgument(0));
|
||||||
|
|
||||||
|
int recordCount = 1500;
|
||||||
|
List<MeteringDataRecordDto> records = new ArrayList<>();
|
||||||
|
LocalDateTime startTime = LocalDateTime.of(2026, 1, 1, 0, 0);
|
||||||
|
|
||||||
|
for (int i = 0; i < recordCount; i++) {
|
||||||
|
LocalDateTime intervalStart = startTime.plusMinutes(i * 15L);
|
||||||
|
LocalDateTime intervalEnd = intervalStart.plusMinutes(15);
|
||||||
|
records.add(new MeteringDataRecordDto(
|
||||||
|
testAtNumber,
|
||||||
|
MeteringDataType.CONSUMPTION,
|
||||||
|
intervalStart,
|
||||||
|
intervalEnd,
|
||||||
|
0.5 + (i % 10) * 0.1,
|
||||||
|
null
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
MeteringDataUploadRequest request = new MeteringDataUploadRequest(
|
||||||
|
DataSource.EMAIL_XLSX,
|
||||||
|
"large_batch.xlsx",
|
||||||
|
records
|
||||||
|
);
|
||||||
|
|
||||||
|
MeteringDataUploadResponse response = meteringDataService.uploadMeteringData(request, UUID.randomUUID());
|
||||||
|
|
||||||
|
assertEquals(UploadStatus.COMPLETED, response.status());
|
||||||
|
assertEquals(recordCount, response.recordCount());
|
||||||
|
assertTrue(response.validationErrors().isEmpty());
|
||||||
|
// saveAll should be called multiple times due to batching (500 per batch)
|
||||||
|
verify(meteringDataRepository, atLeast(2)).saveAll(any());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void uploadMeteringData_duplicateUploadSameDataShouldSucceed() {
|
||||||
|
when(meteringPointRepository.findByAtNumberIn(any())).thenReturn(List.of(testMeteringPoint));
|
||||||
|
when(uploadRepository.save(any())).thenAnswer(invocation -> {
|
||||||
|
MeteringDataUpload upload = invocation.getArgument(0);
|
||||||
|
upload.setId(UUID.randomUUID());
|
||||||
|
return upload;
|
||||||
|
});
|
||||||
|
when(meteringDataRepository.saveAll(any())).thenAnswer(invocation -> invocation.getArgument(0));
|
||||||
|
|
||||||
|
List<MeteringDataRecordDto> records = List.of(
|
||||||
|
new MeteringDataRecordDto(testAtNumber, MeteringDataType.CONSUMPTION,
|
||||||
|
LocalDateTime.of(2026, 1, 1, 0, 0), LocalDateTime.of(2026, 1, 1, 0, 15), 1.25, null),
|
||||||
|
new MeteringDataRecordDto(testAtNumber, MeteringDataType.CONSUMPTION,
|
||||||
|
LocalDateTime.of(2026, 1, 1, 0, 15), LocalDateTime.of(2026, 1, 1, 0, 30), 0.80, null),
|
||||||
|
new MeteringDataRecordDto(testAtNumber, MeteringDataType.CONSUMPTION,
|
||||||
|
LocalDateTime.of(2026, 1, 1, 0, 30), LocalDateTime.of(2026, 1, 1, 0, 45), 1.10, null)
|
||||||
|
);
|
||||||
|
|
||||||
|
MeteringDataUploadRequest request1 = new MeteringDataUploadRequest(
|
||||||
|
DataSource.EMAIL_XLSX,
|
||||||
|
"test.xlsx",
|
||||||
|
records
|
||||||
|
);
|
||||||
|
|
||||||
|
// First upload
|
||||||
|
MeteringDataUploadResponse response1 = meteringDataService.uploadMeteringData(request1, UUID.randomUUID());
|
||||||
|
assertEquals(UploadStatus.COMPLETED, response1.status());
|
||||||
|
assertEquals(3, response1.recordCount());
|
||||||
|
|
||||||
|
// Second upload with same data - should succeed due to deduplication
|
||||||
|
MeteringDataUploadRequest request2 = new MeteringDataUploadRequest(
|
||||||
|
DataSource.EMAIL_XLSX,
|
||||||
|
"test.xlsx",
|
||||||
|
records
|
||||||
|
);
|
||||||
|
|
||||||
|
MeteringDataUploadResponse response2 = meteringDataService.uploadMeteringData(request2, UUID.randomUUID());
|
||||||
|
assertEquals(UploadStatus.COMPLETED, response2.status());
|
||||||
|
assertEquals(3, response2.recordCount());
|
||||||
|
assertTrue(response2.validationErrors().isEmpty());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void uploadMeteringData_largeBatchWithDuplicatesShouldDeduplicate() {
|
||||||
|
when(meteringPointRepository.findByAtNumberIn(any())).thenReturn(List.of(testMeteringPoint));
|
||||||
|
when(uploadRepository.save(any())).thenAnswer(invocation -> {
|
||||||
|
MeteringDataUpload upload = invocation.getArgument(0);
|
||||||
|
upload.setId(UUID.randomUUID());
|
||||||
|
return upload;
|
||||||
|
});
|
||||||
|
when(meteringDataRepository.saveAll(any())).thenAnswer(invocation -> invocation.getArgument(0));
|
||||||
|
|
||||||
|
// Create 2000 records but with only 1000 unique timestamps
|
||||||
|
List<MeteringDataRecordDto> records = new ArrayList<>();
|
||||||
|
LocalDateTime startTime = LocalDateTime.of(2026, 1, 1, 0, 0);
|
||||||
|
|
||||||
|
for (int i = 0; i < 2000; i++) {
|
||||||
|
// Use modulo to create duplicates
|
||||||
|
int index = i % 1000;
|
||||||
|
LocalDateTime intervalStart = startTime.plusMinutes(index * 15L);
|
||||||
|
LocalDateTime intervalEnd = intervalStart.plusMinutes(15);
|
||||||
|
records.add(new MeteringDataRecordDto(
|
||||||
|
testAtNumber,
|
||||||
|
MeteringDataType.CONSUMPTION,
|
||||||
|
intervalStart,
|
||||||
|
intervalEnd,
|
||||||
|
0.5 + (index % 10) * 0.1,
|
||||||
|
null
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
MeteringDataUploadRequest request = new MeteringDataUploadRequest(
|
||||||
|
DataSource.EMAIL_XLSX,
|
||||||
|
"dedup_test.xlsx",
|
||||||
|
records
|
||||||
|
);
|
||||||
|
|
||||||
|
MeteringDataUploadResponse response = meteringDataService.uploadMeteringData(request, UUID.randomUUID());
|
||||||
|
|
||||||
|
assertEquals(UploadStatus.COMPLETED, response.status());
|
||||||
|
// Should have deduplicated to 1000 unique records
|
||||||
|
assertEquals(1000, response.recordCount());
|
||||||
|
assertTrue(response.validationErrors().isEmpty());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user