Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions reporting-pipeline-service/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,7 @@ dependencies {
implementation 'org.springframework.boot:spring-boot-starter-actuator'
implementation 'org.springframework.boot:spring-boot-starter-data-jpa'
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'org.springframework.boot:spring-boot-starter-validation'
developmentOnly 'org.springframework.boot:spring-boot-devtools'

// Kafka
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package gov.cdc.nbs.report.pipeline;

import gov.cdc.nbs.report.pipeline.config.EventProcedureLoggingProperties;
import gov.cdc.nbs.report.pipeline.config.PostProcessingProperties;
import gov.cdc.nbs.report.pipeline.connector.ConnectorProperties;
import gov.cdc.nbs.report.pipeline.lag.LagProperties;
import org.springframework.boot.SpringApplication;
Expand All @@ -15,7 +16,8 @@
@EnableConfigurationProperties({
ConnectorProperties.class,
LagProperties.class,
EventProcedureLoggingProperties.class
EventProcedureLoggingProperties.class,
PostProcessingProperties.class
})
public class ReportingPipelineServiceApplication {

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package gov.cdc.nbs.report.pipeline.config;

import jakarta.validation.constraints.Min;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.validation.annotation.Validated;

@ConfigurationProperties(prefix = "service.post-processing")
@Validated
public record PostProcessingProperties(@Min(0) int maxBatchSize) {}
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import gov.cdc.nbs.report.pipeline.config.PostProcessingProperties;
import gov.cdc.nbs.report.pipeline.postprocessing.repository.InvestigationRepository;
import gov.cdc.nbs.report.pipeline.postprocessing.repository.PostProcRepository;
import gov.cdc.nbs.report.pipeline.postprocessing.repository.model.BackfillData;
Expand Down Expand Up @@ -164,6 +165,7 @@ public class PostProcessingService {

private final ProcessDatamartData dmProcessor;
private final RetryTopicResolver retryTopicResolver;
private final PostProcessingProperties postProcessingProperties;

static final String PAYLOAD = "payload";
static final String SP_EXECUTION_COMPLETED = "Stored proc execution completed: {}";
Expand Down Expand Up @@ -1481,7 +1483,7 @@ protected void processDatamartIds() {
}
}

private String listToParameterString(Collection<Long> inputList) {
private String listToParameterString(Collection<?> inputList) {
return Optional.ofNullable(inputList)
.map(list -> list.stream().map(String::valueOf).distinct().collect(Collectors.joining(",")))
.orElse("");
Expand Down Expand Up @@ -1514,22 +1516,32 @@ private void processTopic(
Collection<Long> ids,
Consumer<String> repositoryMethod,
String... names) {
if (!ids.isEmpty()) {
String idsString = listToParameterString(ids);
String spName = names.length > 0 ? names[0] : entity.getStoredProcedure();
prepareAndLog(keyTopic, idsString, entity.getEntityName(), spName);
String spName = names.length > 0 ? names[0] : entity.getStoredProcedure();
List<List<Long>> chunks =
UidChunker.chunkDistinct(ids, postProcessingProperties.maxBatchSize());
for (int index = 0; index < chunks.size(); index++) {
List<Long> chunk = chunks.get(index);
String idsString = listToParameterString(chunk);
prepareAndLog(
keyTopic, entity.getEntityName(), spName, index + 1, chunks.size(), chunk.size());
repositoryMethod.accept(idsString);
completeLog(spName);
}
}

private void processTopic(
String keyTopic, Entity entity, Collection<String> cds, Consumer<String> repositoryMethod) {
String cdString = cds.stream().distinct().collect(Collectors.joining(","));
String spName = entity.getStoredProcedure();
prepareAndLog(keyTopic, cdString, entity.getEntityName(), spName);
repositoryMethod.accept(cdString);
completeLog(spName);
List<List<String>> chunks =
UidChunker.chunkDistinct(cds, postProcessingProperties.maxBatchSize());
for (int index = 0; index < chunks.size(); index++) {
List<String> chunk = chunks.get(index);
String cdString = listToParameterString(chunk);
prepareAndLog(
keyTopic, entity.getEntityName(), spName, index + 1, chunks.size(), chunk.size());
repositoryMethod.accept(cdString);
completeLog(spName);
}
}

private <T> List<T> processTopic(
Expand All @@ -1540,11 +1552,19 @@ private <T> List<T> processTopic(
Consumer<List<T>> checkResult,
String... names) {
String spName = names.length > 0 ? names[0] : entity.getStoredProcedure();
String idString = listToParameterString(ids);
prepareAndLog(keyTopic, idString, entity.getEntityName(), spName);
List<T> result = repositoryMethod.apply(idString);
checkResult.accept(result);
completeLog(spName);
List<T> result = new ArrayList<>();
List<List<Long>> chunks =
UidChunker.chunkDistinct(ids, postProcessingProperties.maxBatchSize());
for (int index = 0; index < chunks.size(); index++) {
List<Long> chunk = chunks.get(index);
String idString = listToParameterString(chunk);
prepareAndLog(
keyTopic, entity.getEntityName(), spName, index + 1, chunks.size(), chunk.size());
List<T> chunkResult = repositoryMethod.apply(idString);
checkResult.accept(chunkResult);
result.addAll(chunkResult);
completeLog(spName);
}
return result;
}

Expand All @@ -1556,28 +1576,47 @@ private <T> void processTopic(
BiFunction<String, String, List<T>> repositoryMethod,
Consumer<List<T>> checkResult) {
String name = entity.getEntityName();
name = logger.isInfoEnabled() ? StringUtils.capitalize(name) : name;
String idString = listToParameterString(ids);
logger.info(
"Processing {} for topic: {}. Calling stored proc: {} '{}', '{}'",
name,
keyTopic,
entity.getStoredProcedure(),
idString,
vals);
List<T> result = repositoryMethod.apply(idString, vals);
checkResult.accept(result);
completeLog(entity.getStoredProcedure());
String displayName = logger.isInfoEnabled() ? StringUtils.capitalize(name) : name;
List<List<Long>> chunks =
UidChunker.chunkDistinct(ids, postProcessingProperties.maxBatchSize());
for (int index = 0; index < chunks.size(); index++) {
List<Long> chunk = chunks.get(index);
String idString = listToParameterString(chunk);
logger.info(
"Processing {} for topic: {}. Calling stored proc: {} (chunk {}/{}, distinct UIDs: {},"
+ " max batch size: {}, table: '{}')",
displayName,
keyTopic,
entity.getStoredProcedure(),
index + 1,
chunks.size(),
chunk.size(),
postProcessingProperties.maxBatchSize(),
vals);
List<T> result = repositoryMethod.apply(idString, vals);
checkResult.accept(result);
completeLog(entity.getStoredProcedure());
}
}

private void prepareAndLog(String keyTopic, String idString, String name, String spName) {
private void prepareAndLog(
String keyTopic,
String name,
String spName,
int chunkNumber,
int totalChunks,
int distinctCount) {
name = logger.isInfoEnabled() ? StringUtils.capitalize(name) : name;
logger.info(
"Processing {} for topic: {}. Calling stored proc: {} '{}'",
"Processing {} for topic: {}. Calling stored proc: {} (chunk {}/{}, distinct values: {},"
+ " max batch size: {})",
name,
keyTopic,
spName,
idString);
chunkNumber,
totalChunks,
distinctCount,
postProcessingProperties.maxBatchSize());
}

private void completeLog(String sp) {
Expand Down
Loading