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

Reply via email to