Skip to content
Merged
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
191 changes: 183 additions & 8 deletions README.md

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package io.github.nktogo.dataquality.ingestion;

import java.util.UUID;

public interface ValidationRunReportAccess {

ValidationRunResponse getValidationRunForReport(UUID runId);
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
package io.github.nktogo.dataquality.ingestion;

import io.github.nktogo.dataquality.operations.OperationsMetrics;
import io.github.nktogo.dataquality.operations.OperationsMetrics.ValidationProcessingOutcome;
import io.github.nktogo.dataquality.operations.OperationsMetrics.ValidationProcessingSample;
import java.util.List;
import java.util.UUID;
import org.slf4j.Logger;
Expand All @@ -8,38 +11,51 @@
import org.springframework.transaction.annotation.Transactional;

@Service
class ValidationRunService implements ValidationRunAccess {
class ValidationRunService implements ValidationRunAccess, ValidationRunReportAccess {

private static final Logger LOGGER = LoggerFactory.getLogger(ValidationRunService.class);

private final ValidationRunLifecycleService validationRunLifecycleService;
private final ValidationRunRecoveryService validationRunRecoveryService;
private final ValidationRunRepository validationRunRepository;
private final OperationsMetrics operationsMetrics;

ValidationRunService(
ValidationRunLifecycleService validationRunLifecycleService,
ValidationRunRecoveryService validationRunRecoveryService,
ValidationRunRepository validationRunRepository) {
ValidationRunRepository validationRunRepository,
OperationsMetrics operationsMetrics) {
this.validationRunLifecycleService = validationRunLifecycleService;
this.validationRunRecoveryService = validationRunRecoveryService;
this.validationRunRepository = validationRunRepository;
this.operationsMetrics = operationsMetrics;
}

ValidationRunResponse create(UUID fileId, CreateValidationRunRequest request) {
UUID runId = validationRunLifecycleService.createPending(fileId, request.profileId());
operationsMetrics.incrementValidationRunsCreated();
logCreated(runId, fileId, request.profileId());
ValidationProcessingSample processingSample = operationsMetrics.startValidationProcessing();

try {
return validationRunLifecycleService.process(runId);
return recordProcessingResult(validationRunLifecycleService.process(runId), processingSample);
} catch (ValidationProcessingFailureException failure) {
LOGGER.error(
"Validation processing failed for Validation Run '{}'; attempting recovery.",
failure.runId(),
failure.getCause());
logExecutionFailed(runId, fileId, request.profileId(), failure.getCause());
try {
return validationRunRecoveryService.recover(failure);
return recordProcessingResult(
validationRunRecoveryService.recover(failure), processingSample);
} catch (RuntimeException recoveryFailure) {
operationsMetrics.recordValidationProcessing(
processingSample, ValidationProcessingOutcome.ERROR);
logRecoveryFailed(runId, fileId, request.profileId(), recoveryFailure);
recoveryFailure.addSuppressed(failure);
throw recoveryFailure;
}
} catch (RuntimeException executionFailure) {
operationsMetrics.recordValidationProcessing(
processingSample, ValidationProcessingOutcome.ERROR);
logExecutionFailed(runId, fileId, request.profileId(), executionFailure);
throw executionFailure;
}
}

Expand All @@ -52,6 +68,12 @@ List<ValidationRunResponse> getAll() {

@Transactional(readOnly = true)
ValidationRunResponse getById(UUID runId) {
return getValidationRunForReport(runId);
}

@Override
@Transactional(readOnly = true)
public ValidationRunResponse getValidationRunForReport(UUID runId) {
return toResponse(requireExisting(runId));
}

Expand All @@ -67,6 +89,94 @@ private ValidationRun requireExisting(UUID runId) {
.orElseThrow(() -> new ValidationRunNotFoundException(runId));
}

private ValidationRunResponse recordProcessingResult(
ValidationRunResponse response, ValidationProcessingSample processingSample) {
switch (response.status()) {
case COMPLETED -> {
operationsMetrics.recordValidationProcessing(
processingSample, ValidationProcessingOutcome.COMPLETED);
logFinished(response);
}
case FAILED -> {
operationsMetrics.recordValidationProcessing(
processingSample, ValidationProcessingOutcome.FAILED);
logProcessingFailed(response);
}
case PENDING, PROCESSING -> {
operationsMetrics.recordValidationProcessing(
processingSample, ValidationProcessingOutcome.ERROR);
LOGGER
.atError()
.addKeyValue("event", "validation_run.execution_failed")
.addKeyValue("runId", response.id())
.addKeyValue("sourceFileId", response.sourceFileId())
.addKeyValue("profileId", response.profileId())
.addKeyValue("status", response.status())
.log("Validation Run processing returned a nonterminal state.");
}
}
return response;
}

private static void logCreated(UUID runId, UUID fileId, UUID profileId) {
LOGGER
.atInfo()
.addKeyValue("event", "validation_run.created")
.addKeyValue("runId", runId)
.addKeyValue("sourceFileId", fileId)
.addKeyValue("profileId", profileId)
.log("Validation Run created.");
}

private static void logFinished(ValidationRunResponse response) {
addPersistedRunFields(LOGGER.atInfo().addKeyValue("event", "validation_run.finished"), response)
.log("Validation Run processing completed.");
}

private static void logProcessingFailed(ValidationRunResponse response) {
addPersistedRunFields(
LOGGER.atWarn().addKeyValue("event", "validation_run.processing_failed"), response)
.log("Validation Run processing produced a persisted failure.");
}

private static void logExecutionFailed(
UUID runId, UUID fileId, UUID profileId, Throwable failure) {
LOGGER
.atError()
.addKeyValue("event", "validation_run.execution_failed")
.addKeyValue("runId", runId)
.addKeyValue("sourceFileId", fileId)
.addKeyValue("profileId", profileId)
.setCause(failure)
.log("Validation Run execution failed unexpectedly.");
}

private static void logRecoveryFailed(
UUID runId, UUID fileId, UUID profileId, Throwable failure) {
LOGGER
.atError()
.addKeyValue("event", "validation_run.recovery_failed")
.addKeyValue("runId", runId)
.addKeyValue("sourceFileId", fileId)
.addKeyValue("profileId", profileId)
.setCause(failure)
.log("Validation Run failure recovery failed unexpectedly.");
}

private static org.slf4j.spi.LoggingEventBuilder addPersistedRunFields(
org.slf4j.spi.LoggingEventBuilder event, ValidationRunResponse response) {
return event
.addKeyValue("runId", response.id())
.addKeyValue("datasetId", response.datasetId())
.addKeyValue("sourceFileId", response.sourceFileId())
.addKeyValue("profileId", response.profileId())
.addKeyValue("status", response.status())
.addKeyValue("totalRows", response.totalRows())
.addKeyValue("validRows", response.validRows())
.addKeyValue("invalidRows", response.invalidRows())
.addKeyValue("issueCount", response.issueCount());
}

private static ValidationRunResponse toResponse(ValidationRun validationRun) {
return new ValidationRunResponse(
validationRun.getId(),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,178 @@
package io.github.nktogo.dataquality.operations;

import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.Timer;
import java.util.EnumMap;
import java.util.Map;
import java.util.Objects;
import org.springframework.stereotype.Component;

@Component
public final class OperationsMetrics {

private static final String VALIDATION_RUNS_CREATED = "dataquality.validation.runs.created";
private static final String VALIDATION_PROCESSING_DURATION =
"dataquality.validation.processing.duration";
private static final String REPORTS_GENERATED = "dataquality.reports.generated";
private static final String REPORT_GENERATION_DURATION = "dataquality.report.generation.duration";

private final Counter validationRunsCreated;
private final MeterRegistry meterRegistry;
private final Map<ValidationProcessingOutcome, Timer> validationProcessingTimers;
private final Map<ReportFormat, Counter> reportsGenerated;
private final Map<ReportFormat, Map<ReportGenerationOutcome, Timer>> reportGenerationTimers;

OperationsMetrics(MeterRegistry meterRegistry) {
this.meterRegistry = Objects.requireNonNull(meterRegistry, "meterRegistry must not be null");
this.validationRunsCreated =
Counter.builder(VALIDATION_RUNS_CREATED)
.description("Number of Validation Runs created")
.register(meterRegistry);
this.validationProcessingTimers = registerValidationProcessingTimers(meterRegistry);
this.reportsGenerated = registerReportCounters(meterRegistry);
this.reportGenerationTimers = registerReportGenerationTimers(meterRegistry);
}

public void incrementValidationRunsCreated() {
validationRunsCreated.increment();
}

public ValidationProcessingSample startValidationProcessing() {
return new ValidationProcessingSample(Timer.start(meterRegistry));
}

public void recordValidationProcessing(
ValidationProcessingSample sample, ValidationProcessingOutcome outcome) {
Objects.requireNonNull(sample, "sample must not be null")
.stop(
validationProcessingTimers.get(
Objects.requireNonNull(outcome, "outcome must not be null")));
}

public void incrementReportsGenerated(ReportFormat format) {
reportsGenerated.get(Objects.requireNonNull(format, "format must not be null")).increment();
}

public ReportGenerationSample startReportGeneration() {
return new ReportGenerationSample(Timer.start(meterRegistry));
}

public void recordReportGeneration(
ReportGenerationSample sample, ReportFormat format, ReportGenerationOutcome outcome) {
Objects.requireNonNull(sample, "sample must not be null")
.stop(
reportGenerationTimers
.get(Objects.requireNonNull(format, "format must not be null"))
.get(Objects.requireNonNull(outcome, "outcome must not be null")));
}

private static Map<ValidationProcessingOutcome, Timer> registerValidationProcessingTimers(
MeterRegistry meterRegistry) {
Map<ValidationProcessingOutcome, Timer> timers =
new EnumMap<>(ValidationProcessingOutcome.class);
for (ValidationProcessingOutcome outcome : ValidationProcessingOutcome.values()) {
timers.put(
outcome,
Timer.builder(VALIDATION_PROCESSING_DURATION)
.description("Validation Run processing duration")
.tag("outcome", outcome.tagValue)
.register(meterRegistry));
}
return Map.copyOf(timers);
}

private static Map<ReportFormat, Counter> registerReportCounters(MeterRegistry meterRegistry) {
Map<ReportFormat, Counter> counters = new EnumMap<>(ReportFormat.class);
for (ReportFormat format : ReportFormat.values()) {
counters.put(
format,
Counter.builder(REPORTS_GENERATED)
.description("Number of reports generated")
.tag("format", format.tagValue)
.register(meterRegistry));
}
return Map.copyOf(counters);
}

private static Map<ReportFormat, Map<ReportGenerationOutcome, Timer>>
registerReportGenerationTimers(MeterRegistry meterRegistry) {
Map<ReportFormat, Map<ReportGenerationOutcome, Timer>> timers =
new EnumMap<>(ReportFormat.class);
for (ReportFormat format : ReportFormat.values()) {
Map<ReportGenerationOutcome, Timer> formatTimers =
new EnumMap<>(ReportGenerationOutcome.class);
for (ReportGenerationOutcome outcome : ReportGenerationOutcome.values()) {
formatTimers.put(
outcome,
Timer.builder(REPORT_GENERATION_DURATION)
.description("Report generation duration")
.tag("format", format.tagValue)
.tag("outcome", outcome.tagValue)
.register(meterRegistry));
}
timers.put(format, Map.copyOf(formatTimers));
}
return Map.copyOf(timers);
}

public enum ValidationProcessingOutcome {
COMPLETED("completed"),
FAILED("failed"),
ERROR("error");

private final String tagValue;

ValidationProcessingOutcome(String tagValue) {
this.tagValue = tagValue;
}
}

public enum ReportFormat {
JSON("json"),
CSV("csv");

private final String tagValue;

ReportFormat(String tagValue) {
this.tagValue = tagValue;
}
}

public enum ReportGenerationOutcome {
SUCCESS("success"),
ERROR("error");

private final String tagValue;

ReportGenerationOutcome(String tagValue) {
this.tagValue = tagValue;
}
}

public static final class ValidationProcessingSample {

private final Timer.Sample sample;

private ValidationProcessingSample(Timer.Sample sample) {
this.sample = sample;
}

private void stop(Timer timer) {
sample.stop(timer);
}
}

public static final class ReportGenerationSample {

private final Timer.Sample sample;

private ReportGenerationSample(Timer.Sample sample) {
this.sample = sample;
}

private void stop(Timer timer) {
sample.stop(timer);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package io.github.nktogo.dataquality.reporting;

final class InvalidReportFormatException extends RuntimeException {

InvalidReportFormatException() {
super("Query parameter 'format' must be exactly one of: json, csv.");
}
}
Loading