This is an automated email from the ASF dual-hosted git repository.
exceptionfactory pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new c3e7aa25433 NIFI-16161 Added the ability to validate JSON schemas
before use in ValidateJson (#11503)
c3e7aa25433 is described below
commit c3e7aa25433fb281cae3077257a19d45e95fb588
Author: dan-s1 <[email protected]>
AuthorDate: Mon Aug 3 17:00:30 2026 -0400
NIFI-16161 Added the ability to validate JSON schemas before use in
ValidateJson (#11503)
Signed-off-by: David Handermann <[email protected]>
---
.../nifi/processors/standard/ValidateJson.java | 42 +++++++++++++++++-----
.../nifi/processors/standard/TestValidateJson.java | 19 +++++-----
2 files changed, 43 insertions(+), 18 deletions(-)
diff --git
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ValidateJson.java
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ValidateJson.java
index 83dd2e7c836..af7caf451cd 100644
---
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ValidateJson.java
+++
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ValidateJson.java
@@ -61,6 +61,7 @@ import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.io.LineNumberReader;
+import java.io.UncheckedIOException;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
@@ -286,21 +287,34 @@ public class ValidateJson extends AbstractProcessor {
.build());
}
+ if
(schemaAccessStrategy.equals(JsonSchemaStrategy.SCHEMA_CONTENT_PROPERTY) &&
validationContext.getProperty(SCHEMA_CONTENT).isSet()) {
+ try {
+ readSchema(validationContext);
+ } catch (final Exception e) {
+ final String reason;
+ final Throwable cause = e.getCause();
+ if (cause == null) {
+ reason = e.getMessage();
+ } else {
+ reason = "%s [%s]".formatted(cause.getMessage(),
e.getMessage());
+ }
+
+ final String message = "JSON schema not valid:
%s".formatted(reason);
+ validationResults.add(new ValidationResult.Builder()
+ .valid(false)
+ .subject(SCHEMA_CONTENT.getDisplayName())
+ .explanation(message)
+ .build());
+ }
+ }
return validationResults;
}
@OnScheduled
public void onScheduled(final ProcessContext context) throws IOException {
switch (getSchemaAccessStrategy(context)) {
- case SCHEMA_NAME_PROPERTY ->
- jsonSchemaRegistry =
context.getProperty(SCHEMA_REGISTRY).asControllerService(JsonSchemaRegistry.class);
- case SCHEMA_CONTENT_PROPERTY -> {
- try (final InputStream inputStream =
context.getProperty(SCHEMA_CONTENT).asResource().read()) {
- final SchemaVersion schemaVersion =
SchemaVersion.valueOf(context.getProperty(SCHEMA_VERSION).getValue());
- final SchemaRegistry registry =
schemaRegistries.get(schemaVersion);
- schema = registry.getSchema(inputStream);
- }
- }
+ case SCHEMA_NAME_PROPERTY -> jsonSchemaRegistry =
context.getProperty(SCHEMA_REGISTRY).asControllerService(JsonSchemaRegistry.class);
+ case SCHEMA_CONTENT_PROPERTY -> schema = readSchema(context);
}
final int maxStringLength =
context.getProperty(MAX_STRING_LENGTH).asDataSize(DataUnit.B).intValue();
@@ -339,6 +353,16 @@ public class ValidateJson extends AbstractProcessor {
}
}
+ private Schema readSchema(final PropertyContext context) {
+ try (final InputStream inputStream =
context.getProperty(SCHEMA_CONTENT).asResource().read()) {
+ final SchemaVersion schemaVersion =
SchemaVersion.valueOf(context.getProperty(SCHEMA_VERSION).getValue());
+ final SchemaRegistry registry =
schemaRegistries.get(schemaVersion);
+ return registry.getSchema(inputStream);
+ } catch (final IOException ioe) {
+ throw new UncheckedIOException("Read JSON schema failed", ioe);
+ }
+ }
+
private void validateFlowFile(final ProcessSession session, final FlowFile
flowFile) {
final Schema currentSchema = schema;
diff --git
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestValidateJson.java
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestValidateJson.java
index e5d1cbc7764..19ead4e9eca 100644
---
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestValidateJson.java
+++
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestValidateJson.java
@@ -16,7 +16,7 @@
*/
package org.apache.nifi.processors.standard;
-import org.apache.commons.lang3.exception.ExceptionUtils;
+import org.apache.nifi.components.ValidationResult;
import org.apache.nifi.controller.AbstractControllerService;
import org.apache.nifi.json.schema.JsonSchema;
import org.apache.nifi.json.schema.SchemaVersion;
@@ -40,6 +40,7 @@ import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
+import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.stream.Stream;
@@ -259,11 +260,11 @@ class TestValidateJson {
runner.setProperty(ValidateJson.SCHEMA_CONTENT, schema);
runner.setProperty(JsonSchemaRegistryComponent.SCHEMA_VERSION,
SCHEMA_VERSION);
runner.enqueue(getFileContent("simple-example-with-comments.json"));
- runner.assertValid();
- final AssertionFailedError assertionFailedError =
assertThrows(AssertionFailedError.class, () -> runner.run());
- final String stackTrace =
ExceptionUtils.getStackTrace(assertionFailedError);
- assertTrue(stackTrace.contains("JsonParseException") &&
!stackTrace.contains("FileNotFoundException"));
+ runner.assertNotValid();
+ final Collection<ValidationResult> validationResults =
runner.validate();
+ final String explanation =
validationResults.iterator().next().getExplanation();
+ assertTrue(explanation.contains("JsonParseException") &&
!explanation.contains("FileNotFoundException"));
}
@ParameterizedTest
@@ -272,11 +273,11 @@ class TestValidateJson {
runner.setProperty(ValidateJson.SCHEMA_CONTENT, schema);
runner.setProperty(JsonSchemaRegistryComponent.SCHEMA_VERSION,
SCHEMA_VERSION);
runner.enqueue(getFileContent("simple-example-with-comments.json"));
- runner.assertValid();
- final AssertionFailedError assertionFailedError =
assertThrows(AssertionFailedError.class, () -> runner.run());
- final String stackTrace =
ExceptionUtils.getStackTrace(assertionFailedError);
- assertTrue(stackTrace.contains("JsonParseException"));
+ runner.assertNotValid();
+ final Collection<ValidationResult> validationResults =
runner.validate();
+ final String explanation =
validationResults.iterator().next().getExplanation();
+ assertTrue(explanation.contains("JsonParseException"));
}
@ParameterizedTest