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 91cc43a62c7 NIFI-15587 Added extensible FieldValidator and 
RecordValidator support to the Record API (#10890)
91cc43a62c7 is described below

commit 91cc43a62c726d11d6aedbfe16de961d786d5f18
Author: Pierre Villard <[email protected]>
AuthorDate: Mon Aug 31 20:17:26 2026 +0200

    NIFI-15587 Added extensible FieldValidator and RecordValidator support to 
the Record API (#10890)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../record/validation/DefaultValidationError.java  | 131 +++++
 .../record/validation/FieldValidator.java          |  40 ++
 .../record/validation/RecordValidator.java         |  43 ++
 .../record/validation/SchemaValidators.java        |  75 +++
 .../validation/TestDefaultValidationError.java     | 147 +++++
 .../schema/validation/SchemaValidationContext.java |  11 +
 .../schema/validation/StandardSchemaValidator.java | 128 ++++-
 .../validation/TestStandardSchemaValidator.java    | 595 +++++++++++++++++++++
 8 files changed, 1142 insertions(+), 28 deletions(-)

diff --git 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/DefaultValidationError.java
 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/DefaultValidationError.java
new file mode 100644
index 00000000000..cbf6c85c465
--- /dev/null
+++ 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/DefaultValidationError.java
@@ -0,0 +1,131 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.serialization.record.validation;
+
+import java.util.Objects;
+import java.util.Optional;
+import java.util.StringJoiner;
+
+/**
+ * Basic implementation of {@link ValidationError} that can be used by 
validators in modules that
+ * cannot depend on higher level utility classes. Instances are immutable and 
thread-safe.
+ */
+public class DefaultValidationError implements ValidationError {
+    private final Optional<String> fieldName;
+    private final Optional<Object> inputValue;
+    private final String explanation;
+    private final ValidationErrorType type;
+
+    private DefaultValidationError(final Builder builder) {
+        this.fieldName = Optional.ofNullable(builder.fieldName);
+        this.inputValue = Optional.ofNullable(builder.inputValue);
+        this.explanation = Objects.requireNonNull(builder.explanation, 
"Explanation is required");
+        this.type = Objects.requireNonNull(builder.type, "Validation error 
type is required");
+    }
+
+    @Override
+    public Optional<String> getFieldName() {
+        return fieldName;
+    }
+
+    @Override
+    public Optional<Object> getInputValue() {
+        return inputValue;
+    }
+
+    @Override
+    public String getExplanation() {
+        return explanation;
+    }
+
+    @Override
+    public ValidationErrorType getType() {
+        return type;
+    }
+
+    @Override
+    public String toString() {
+        final StringJoiner joiner = new StringJoiner(", ", 
"DefaultValidationError[", "]");
+        fieldName.ifPresent(name -> joiner.add("field=" + name));
+        inputValue.ifPresent(value -> joiner.add("value=" + value));
+        joiner.add("type=" + type);
+        joiner.add("explanation=" + explanation);
+        return joiner.toString();
+    }
+
+    @Override
+    public int hashCode() {
+        return 31 + 17 * fieldName.hashCode() + 17 * inputValue.hashCode() + 
17 * explanation.hashCode();
+    }
+
+    @Override
+    public boolean equals(final Object obj) {
+        if (obj == this) {
+            return true;
+        }
+        if (obj == null) {
+            return false;
+        }
+        if (!(obj instanceof ValidationError other)) {
+            return false;
+        }
+        return getFieldName().equals(other.getFieldName()) && 
getInputValue().equals(other.getInputValue()) && 
getExplanation().equals(other.getExplanation());
+    }
+
+    /**
+     * Creates a builder for constructing immutable {@link 
DefaultValidationError} instances.
+     *
+     * @return builder instance
+     */
+    public static Builder builder() {
+        return new Builder();
+    }
+
+    public static final class Builder {
+        private String fieldName;
+        private Object inputValue;
+        private String explanation;
+        private ValidationErrorType type = ValidationErrorType.INVALID_FIELD;
+
+        private Builder() {
+        }
+
+        public Builder fieldName(final String fieldName) {
+            this.fieldName = fieldName;
+            return this;
+        }
+
+        public Builder inputValue(final Object inputValue) {
+            this.inputValue = inputValue;
+            return this;
+        }
+
+        public Builder explanation(final String explanation) {
+            this.explanation = explanation;
+            return this;
+        }
+
+        public Builder type(final ValidationErrorType type) {
+            this.type = type;
+            return this;
+        }
+
+        public DefaultValidationError build() {
+            return new DefaultValidationError(this);
+        }
+    }
+}
diff --git 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/FieldValidator.java
 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/FieldValidator.java
new file mode 100644
index 00000000000..8ec0ef0e550
--- /dev/null
+++ 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/FieldValidator.java
@@ -0,0 +1,40 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.serialization.record.validation;
+
+import java.util.Collection;
+
+/**
+ * Provides validation for an individual record field value. Field Validators 
are expected to be immutable and thread-safe.
+ * If a validator needs field metadata (such as the data type), it should 
capture that information at construction time.
+ */
+public interface FieldValidator {
+
+    /**
+     * Validates the provided value for a field at the given path.
+     *
+     * @param fieldPath the path of the field being validated (used for clear 
diagnostics)
+     * @param value the value of the field for the record currently being 
validated
+     * @return a collection of validation errors. The collection must be empty 
when the value is valid.
+     */
+    Collection<ValidationError> validate(String fieldPath, Object value);
+
+    /**
+     * @return a short human readable description of what the validator 
enforces
+     */
+    String getDescription();
+}
diff --git 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/RecordValidator.java
 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/RecordValidator.java
new file mode 100644
index 00000000000..5da522d6a94
--- /dev/null
+++ 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/RecordValidator.java
@@ -0,0 +1,43 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.serialization.record.validation;
+
+import org.apache.nifi.serialization.record.Record;
+
+import java.util.Collection;
+
+/**
+ * Provides validation logic for an entire {@link Record} instance.
+ * Record Validators are expected to be immutable and thread-safe.
+ * If a validator needs schema metadata, it should capture that information at 
construction time.
+ */
+public interface RecordValidator {
+
+    /**
+     * Validates the provided record.
+     *
+     * @param record the record instance to validate
+     * @param fieldPath the path within the overall document that identifies 
the record (root records use the empty string)
+     * @return a collection of validation errors. The collection must be empty 
when the record is valid.
+     */
+    Collection<ValidationError> validate(Record record, String fieldPath);
+
+    /**
+     * @return a short human readable description of what the validator 
enforces
+     */
+    String getDescription();
+}
diff --git 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/SchemaValidators.java
 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/SchemaValidators.java
new file mode 100644
index 00000000000..9ce3f56a950
--- /dev/null
+++ 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/validation/SchemaValidators.java
@@ -0,0 +1,75 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.serialization.record.validation;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Immutable container that carries field-level and record-level validators 
independently of the schema model.
+ * Field validators are keyed by field name. Record validators apply to the 
entire record.
+ * Nested validators mirror the schema tree: a RECORD-type or ARRAY-of-RECORD 
field maps to a child
+ * {@code SchemaValidators} that applies to the nested record.
+ */
+public class SchemaValidators {
+    public static final SchemaValidators EMPTY = new 
SchemaValidators(Map.of(), List.of(), Map.of());
+
+    private final Map<String, List<FieldValidator>> fieldValidators;
+    private final List<RecordValidator> recordValidators;
+    private final Map<String, SchemaValidators> nestedValidators;
+
+    public SchemaValidators(final Map<String, List<FieldValidator>> 
fieldValidators, final List<RecordValidator> recordValidators) {
+        this(fieldValidators, recordValidators, Map.of());
+    }
+
+    public SchemaValidators(final Map<String, List<FieldValidator>> 
fieldValidators, final List<RecordValidator> recordValidators,
+            final Map<String, SchemaValidators> nestedValidators) {
+        if (fieldValidators == null || fieldValidators.isEmpty()) {
+            this.fieldValidators = Map.of();
+        } else {
+            final Map<String, List<FieldValidator>> defensiveCopy = new 
HashMap<>(fieldValidators.size());
+            for (final Map.Entry<String, List<FieldValidator>> entry : 
fieldValidators.entrySet()) {
+                defensiveCopy.put(entry.getKey(), 
List.copyOf(entry.getValue()));
+            }
+            this.fieldValidators = Collections.unmodifiableMap(defensiveCopy);
+        }
+        this.recordValidators = recordValidators == null || 
recordValidators.isEmpty() ? List.of() : List.copyOf(recordValidators);
+        this.nestedValidators = nestedValidators == null || 
nestedValidators.isEmpty() ? Map.of() : Map.copyOf(nestedValidators);
+    }
+
+    public List<FieldValidator> getFieldValidators(final String fieldName) {
+        return fieldValidators.getOrDefault(fieldName, List.of());
+    }
+
+    public Map<String, List<FieldValidator>> getAllFieldValidators() {
+        return fieldValidators;
+    }
+
+    public List<RecordValidator> getRecordValidators() {
+        return recordValidators;
+    }
+
+    public SchemaValidators getNestedValidators(final String fieldName) {
+        return nestedValidators.getOrDefault(fieldName, EMPTY);
+    }
+
+    public boolean isEmpty() {
+        return fieldValidators.isEmpty() && recordValidators.isEmpty() && 
nestedValidators.isEmpty();
+    }
+}
diff --git 
a/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/validation/TestDefaultValidationError.java
 
b/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/validation/TestDefaultValidationError.java
new file mode 100644
index 00000000000..ab55655305d
--- /dev/null
+++ 
b/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/validation/TestDefaultValidationError.java
@@ -0,0 +1,147 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.serialization.record.validation;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class TestDefaultValidationError {
+
+    @Test
+    void testBuilderProducesCorrectFields() {
+        final DefaultValidationError error = DefaultValidationError.builder()
+                .fieldName("/root/name")
+                .inputValue("badValue")
+                .type(ValidationErrorType.MISSING_FIELD)
+                .explanation("field is missing")
+                .build();
+
+        assertEquals(Optional.of("/root/name"), error.getFieldName());
+        assertEquals(Optional.of("badValue"), error.getInputValue());
+        assertEquals(ValidationErrorType.MISSING_FIELD, error.getType());
+        assertEquals("field is missing", error.getExplanation());
+    }
+
+    @Test
+    void testExplanationIsRequired() {
+        final DefaultValidationError.Builder builder = 
DefaultValidationError.builder()
+                .fieldName("/field")
+                .type(ValidationErrorType.INVALID_FIELD);
+
+        assertThrows(NullPointerException.class, builder::build);
+    }
+
+    @Test
+    void testTypeIsRequired() {
+        final DefaultValidationError.Builder builder = 
DefaultValidationError.builder()
+                .explanation("some explanation")
+                .type(null);
+
+        assertThrows(NullPointerException.class, builder::build);
+    }
+
+    @Test
+    void testDefaultTypeIsInvalidField() {
+        final DefaultValidationError error = DefaultValidationError.builder()
+                .explanation("some explanation")
+                .build();
+
+        assertEquals(ValidationErrorType.INVALID_FIELD, error.getType());
+    }
+
+    @Test
+    void testOptionalFieldsDefaultToEmpty() {
+        final DefaultValidationError error = DefaultValidationError.builder()
+                .explanation("explanation only")
+                .build();
+
+        assertEquals(Optional.empty(), error.getFieldName());
+        assertEquals(Optional.empty(), error.getInputValue());
+    }
+
+    @Test
+    void testEqualsAndHashCode() {
+        final DefaultValidationError error1 = DefaultValidationError.builder()
+                .fieldName("/field")
+                .inputValue(42)
+                .type(ValidationErrorType.INVALID_FIELD)
+                .explanation("bad value")
+                .build();
+
+        final DefaultValidationError error2 = DefaultValidationError.builder()
+                .fieldName("/field")
+                .inputValue(42)
+                .type(ValidationErrorType.MISSING_FIELD)
+                .explanation("bad value")
+                .build();
+
+        assertEquals(error1, error2);
+        assertEquals(error1.hashCode(), error2.hashCode());
+    }
+
+    @Test
+    void testNotEqualWhenExplanationDiffers() {
+        final DefaultValidationError error1 = DefaultValidationError.builder()
+                .fieldName("/field")
+                .explanation("explanation A")
+                .build();
+
+        final DefaultValidationError error2 = DefaultValidationError.builder()
+                .fieldName("/field")
+                .explanation("explanation B")
+                .build();
+
+        assertNotEquals(error1, error2);
+    }
+
+    @Test
+    void testNotEqualWhenFieldNameDiffers() {
+        final DefaultValidationError error1 = DefaultValidationError.builder()
+                .fieldName("/fieldA")
+                .explanation("same")
+                .build();
+
+        final DefaultValidationError error2 = DefaultValidationError.builder()
+                .fieldName("/fieldB")
+                .explanation("same")
+                .build();
+
+        assertNotEquals(error1, error2);
+    }
+
+    @Test
+    void testToStringContainsAllFields() {
+        final DefaultValidationError error = DefaultValidationError.builder()
+                .fieldName("/root/name")
+                .inputValue("badValue")
+                .type(ValidationErrorType.MISSING_FIELD)
+                .explanation("field is missing")
+                .build();
+
+        final String result = error.toString();
+        assertTrue(result.contains("field=/root/name"));
+        assertTrue(result.contains("value=badValue"));
+        assertTrue(result.contains("type=MISSING_FIELD"));
+        assertTrue(result.contains("explanation=field is missing"));
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/SchemaValidationContext.java
 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/SchemaValidationContext.java
index 53711de29b6..55c53d01a15 100644
--- 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/SchemaValidationContext.java
+++ 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/SchemaValidationContext.java
@@ -18,16 +18,23 @@
 package org.apache.nifi.schema.validation;
 
 import org.apache.nifi.serialization.record.RecordSchema;
+import org.apache.nifi.serialization.record.validation.SchemaValidators;
 
 public class SchemaValidationContext {
     private final RecordSchema schema;
     private final boolean allowExtraFields;
     private final boolean strictTypeChecking;
+    private final SchemaValidators schemaValidators;
 
     public SchemaValidationContext(final RecordSchema schema, final boolean 
allowExtraFields, final boolean strictTypeChecking) {
+        this(schema, allowExtraFields, strictTypeChecking, 
SchemaValidators.EMPTY);
+    }
+
+    public SchemaValidationContext(final RecordSchema schema, final boolean 
allowExtraFields, final boolean strictTypeChecking, final SchemaValidators 
schemaValidators) {
         this.schema = schema;
         this.allowExtraFields = allowExtraFields;
         this.strictTypeChecking = strictTypeChecking;
+        this.schemaValidators = schemaValidators == null ? 
SchemaValidators.EMPTY : schemaValidators;
     }
 
     public RecordSchema getSchema() {
@@ -41,4 +48,8 @@ public class SchemaValidationContext {
     public boolean isStrictTypeChecking() {
         return strictTypeChecking;
     }
+
+    public SchemaValidators getSchemaValidators() {
+        return schemaValidators;
+    }
 }
diff --git 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/StandardSchemaValidator.java
 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/StandardSchemaValidator.java
index bf96a30da86..b9a8808ec26 100644
--- 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/StandardSchemaValidator.java
+++ 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/main/java/org/apache/nifi/schema/validation/StandardSchemaValidator.java
@@ -28,8 +28,11 @@ import 
org.apache.nifi.serialization.record.type.EnumDataType;
 import org.apache.nifi.serialization.record.type.MapDataType;
 import org.apache.nifi.serialization.record.type.RecordDataType;
 import org.apache.nifi.serialization.record.util.DataTypeUtils;
+import org.apache.nifi.serialization.record.validation.FieldValidator;
 import org.apache.nifi.serialization.record.validation.RecordSchemaValidator;
+import org.apache.nifi.serialization.record.validation.RecordValidator;
 import org.apache.nifi.serialization.record.validation.SchemaValidationResult;
+import org.apache.nifi.serialization.record.validation.SchemaValidators;
 import org.apache.nifi.serialization.record.validation.ValidationError;
 import org.apache.nifi.serialization.record.validation.ValidationErrorType;
 
@@ -37,6 +40,7 @@ import java.math.BigInteger;
 import java.sql.Date;
 import java.sql.Time;
 import java.sql.Timestamp;
+import java.util.Collection;
 import java.util.List;
 import java.util.Map;
 
@@ -49,10 +53,10 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
 
     @Override
     public SchemaValidationResult validate(final Record record) {
-        return validate(record, validationContext.getSchema(), "");
+        return validate(record, validationContext.getSchema(), "", 
validationContext.getSchemaValidators());
     }
 
-    private SchemaValidationResult validate(final Record record, final 
RecordSchema schema, final String fieldPrefix) {
+    private SchemaValidationResult validate(final Record record, final 
RecordSchema schema, final String fieldPrefix, final SchemaValidators 
validators) {
         // Ensure that for every field in the schema, the type is correct (if 
we care) and that
         // a value is present (unless it is nullable).
         final StandardSchemaValidationResult result = new 
StandardSchemaValidationResult();
@@ -96,8 +100,11 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
                 continue;
             }
 
+            final String fieldPath = concat(fieldPrefix, field);
+            applyFieldValidators(field, fieldPath, rawValue, result, 
validators);
+
             // Now that we have the 'canonical data type', we check if it is a 
Record. If so, we need to validate each sub-field.
-            verifyComplexType(dataType, rawValue, result, fieldPrefix, field);
+            verifyComplexType(canonicalDataType, rawValue, result, 
fieldPrefix, field, validators);
         }
 
         if (!validationContext.isExtraFieldAllowed()) {
@@ -108,20 +115,51 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
             }
         }
 
+        applyRecordValidators(record, fieldPrefix, result, validators);
+
         return result;
     }
 
-    private void verifyComplexType(final DataType dataType, final Object 
rawValue, final StandardSchemaValidationResult result, final String 
fieldPrefix, final RecordField field) {
-        // If the field type is RECORD, or if the field type is a CHOICE that 
allows for a RECORD and the value is a RECORD, then we
-        // need to dig into each of the sub-fields. To do this, we first need 
to determine the 'canonical data type'.
-        final DataType canonicalDataType = getCanonicalDataType(dataType, 
rawValue, result, fieldPrefix, field);
-        if (canonicalDataType == null) {
+    private void applyFieldValidators(final RecordField field, final String 
fieldPath, final Object value, final StandardSchemaValidationResult result,
+            final SchemaValidators validators) {
+        if (value == null) {
             return;
         }
 
-        // Now that we have the 'canonical data type', we check if it is a 
Record. If so, we need to validate each sub-field.
+        for (final FieldValidator validator : 
validators.getFieldValidators(field.getFieldName())) {
+            final Collection<ValidationError> errors = 
validator.validate(fieldPath, value);
+            if (errors == null || errors.isEmpty()) {
+                continue;
+            }
+
+            for (final ValidationError validationError : errors) {
+                result.addValidationError(validationError);
+            }
+        }
+    }
+
+    private void applyRecordValidators(final Record record, final String 
fieldPath, final StandardSchemaValidationResult result, final SchemaValidators 
validators) {
+        final List<RecordValidator> recordValidators = 
validators.getRecordValidators();
+        if (recordValidators.isEmpty()) {
+            return;
+        }
+
+        for (final RecordValidator recordValidator : recordValidators) {
+            final Collection<ValidationError> validationErrors = 
recordValidator.validate(record, fieldPath);
+            if (validationErrors == null || validationErrors.isEmpty()) {
+                continue;
+            }
+
+            for (final ValidationError validationError : validationErrors) {
+                result.addValidationError(validationError);
+            }
+        }
+    }
+
+    private void verifyComplexType(final DataType canonicalDataType, final 
Object rawValue, final StandardSchemaValidationResult result, final String 
fieldPrefix,
+            final RecordField field, final SchemaValidators validators) {
         if (canonicalDataType.getFieldType() == RecordFieldType.RECORD) {
-            verifyChildRecord(canonicalDataType, rawValue, dataType, result, 
field, fieldPrefix);
+            verifyChildRecord(canonicalDataType, rawValue, result, field, 
fieldPrefix, validators);
         }
 
         if (canonicalDataType.getFieldType() == RecordFieldType.ARRAY) {
@@ -129,9 +167,11 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
             final DataType elementType = arrayDataType.getElementType();
             final Object[] arrayObject = (Object[]) rawValue;
 
+            final String arrayPath = concat(fieldPrefix, field);
+            final SchemaValidators childValidators = 
validators.getNestedValidators(field.getFieldName());
             int i = 0;
             for (final Object arrayValue : arrayObject) {
-                verifyComplexType(elementType, arrayValue, result, fieldPrefix 
+ "[" + i + "]", field);
+                validateArrayElement(elementType, arrayValue, result, 
arrayPath + "[" + i + "]", childValidators);
                 i++;
             }
         }
@@ -156,28 +196,60 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
         return canonicalDataType;
     }
 
-    private void verifyChildRecord(final DataType canonicalDataType, final 
Object rawValue, final DataType expectedDataType, final 
StandardSchemaValidationResult result,
-        final RecordField field, final String fieldPrefix) {
-        // Now that we have the 'canonical data type', we check if it is a 
Record. If so, we need to validate each sub-field.
-        if (canonicalDataType.getFieldType() == RecordFieldType.RECORD) {
-            if (!(rawValue instanceof Record)) { // sanity check
-                result.addValidationError(new 
StandardValidationError(concat(fieldPrefix, field), rawValue, 
ValidationErrorType.INVALID_FIELD,
-                    "Value is of type " + classNameOrNull(rawValue) + " but 
was expected to be of type " + expectedDataType));
+    private void verifyChildRecord(final DataType canonicalDataType, final 
Object rawValue, final StandardSchemaValidationResult result,
+            final RecordField field, final String fieldPrefix, final 
SchemaValidators validators) {
+        if (!(rawValue instanceof Record)) {
+            result.addValidationError(new 
StandardValidationError(concat(fieldPrefix, field), rawValue, 
ValidationErrorType.INVALID_FIELD,
+                "Value is of type " + classNameOrNull(rawValue) + " but was 
expected to be of type " + canonicalDataType));
+            return;
+        }
 
-                return;
-            }
+        final RecordDataType recordDataType = (RecordDataType) 
canonicalDataType;
+        final RecordSchema childSchema = recordDataType.getChildSchema();
+
+        final String fullChildFieldName = concat(fieldPrefix, field);
+        final SchemaValidators childValidators = 
validators.getNestedValidators(field.getFieldName());
+        final SchemaValidationResult childValidationResult = validate((Record) 
rawValue, childSchema, fullChildFieldName, childValidators);
+        if (childValidationResult.isValid()) {
+            return;
+        }
+
+        for (final ValidationError validationError : 
childValidationResult.getValidationErrors()) {
+            result.addValidationError(validationError);
+        }
+    }
 
-            final RecordDataType recordDataType = (RecordDataType) 
canonicalDataType;
-            final RecordSchema childSchema = recordDataType.getChildSchema();
+    private void validateArrayElement(final DataType elementType, final Object 
rawValue, final StandardSchemaValidationResult result,
+            final String elementPath, final SchemaValidators validators) {
+        if (rawValue == null) {
+            return;
+        }
 
-            final String fullChildFieldName = concat(fieldPrefix, field);
-            final SchemaValidationResult childValidationResult = 
validate((Record) rawValue, childSchema, fullChildFieldName);
-            if (childValidationResult.isValid()) {
+        final DataType resolvedType;
+        if (elementType.getFieldType() == RecordFieldType.CHOICE) {
+            resolvedType = DataTypeUtils.chooseDataType(rawValue, 
(ChoiceDataType) elementType);
+            if (resolvedType == null) {
+                result.addValidationError(new 
StandardValidationError(elementPath, rawValue, 
ValidationErrorType.INVALID_FIELD,
+                    "Value is of type " + classNameOrNull(rawValue) + " but 
was expected to be of type " + elementType));
                 return;
             }
+        } else {
+            resolvedType = elementType;
+        }
 
-            for (final ValidationError validationError : 
childValidationResult.getValidationErrors()) {
-                result.addValidationError(validationError);
+        if (resolvedType.getFieldType() == RecordFieldType.RECORD && rawValue 
instanceof Record) {
+            final RecordDataType recordDataType = (RecordDataType) 
resolvedType;
+            final SchemaValidationResult childResult = validate((Record) 
rawValue, recordDataType.getChildSchema(), elementPath, validators);
+            for (final ValidationError error : 
childResult.getValidationErrors()) {
+                result.addValidationError(error);
+            }
+        } else if (resolvedType.getFieldType() == RecordFieldType.ARRAY && 
rawValue instanceof Object[]) {
+            final ArrayDataType nestedArrayType = (ArrayDataType) resolvedType;
+            final Object[] nestedArray = (Object[]) rawValue;
+            int i = 0;
+            for (final Object nestedValue : nestedArray) {
+                validateArrayElement(nestedArrayType.getElementType(), 
nestedValue, result, elementPath + "[" + i + "]", validators);
+                i++;
             }
         }
     }
@@ -303,7 +375,7 @@ public class StandardSchemaValidator implements 
RecordSchemaValidator {
         return fieldPrefix + "/" + field.getFieldName();
     }
 
-    private String classNameOrNull(Object value) {
+    private String classNameOrNull(final Object value) {
         return value == null ? "null" : value.getClass().getName();
     }
 }
diff --git 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/test/java/org/apache/nifi/schema/validation/TestStandardSchemaValidator.java
 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/test/java/org/apache/nifi/schema/validation/TestStandardSchemaValidator.java
index 6971cd1b1dd..fd65727fac1 100644
--- 
a/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/test/java/org/apache/nifi/schema/validation/TestStandardSchemaValidator.java
+++ 
b/nifi-extension-bundles/nifi-extension-utils/nifi-record-utils/nifi-standard-record-utils/src/test/java/org/apache/nifi/schema/validation/TestStandardSchemaValidator.java
@@ -25,7 +25,11 @@ import org.apache.nifi.serialization.record.RecordField;
 import org.apache.nifi.serialization.record.RecordFieldType;
 import org.apache.nifi.serialization.record.RecordSchema;
 import org.apache.nifi.serialization.record.type.EnumDataType;
+import org.apache.nifi.serialization.record.validation.DefaultValidationError;
+import org.apache.nifi.serialization.record.validation.FieldValidator;
+import org.apache.nifi.serialization.record.validation.RecordValidator;
 import org.apache.nifi.serialization.record.validation.SchemaValidationResult;
+import org.apache.nifi.serialization.record.validation.SchemaValidators;
 import org.apache.nifi.serialization.record.validation.ValidationError;
 import org.apache.nifi.serialization.record.validation.ValidationErrorType;
 import org.junit.jupiter.api.Test;
@@ -45,6 +49,7 @@ import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -59,6 +64,13 @@ public class TestStandardSchemaValidator {
     private static final Long MAX_PRECISE_WHOLE_IN_FLOAT = 
Double.valueOf(Math.pow(2, FLOAT_BITS_PRECISION)).longValue();
     private static final Long MAX_PRECISE_WHOLE_IN_DOUBLE = 
Double.valueOf(Math.pow(2, DOUBLE_BITS_PRECISION)).longValue();
 
+    private static final String FIELD_VALUE = "value";
+    private static final String FIELD_NAME = "name";
+    private static final String FIELD_CHILD = "child";
+    private static final String FIELD_ID = "id";
+    private static final String FIELD_KEY = "key";
+    private static final String FIELD_ITEMS = "items";
+
     private static final Set<RecordFieldType> NUMERIC_TYPES = new 
HashSet<>(Arrays.asList(
             RecordFieldType.BYTE,
             RecordFieldType.SHORT,
@@ -433,6 +445,589 @@ public class TestStandardSchemaValidator {
         assertTrue(result.getValidationErrors().isEmpty());
     }
 
+    @Test
+    public void testFieldValidatorIsApplied() {
+        final AtomicBoolean invoked = new AtomicBoolean(false);
+        final FieldValidator fieldValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                invoked.set(true);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("value must be 'pass'")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "validator for testFieldValidatorIsApplied";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType(), null, Collections.emptySet(), false);
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"fail"));
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(FIELD_VALUE, List.of(fieldValidator)), List.of());
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, false, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertTrue(invoked.get());
+        assertFalse(result.isValid());
+        final ValidationError validationError = 
result.getValidationErrors().iterator().next();
+        assertEquals(ValidationErrorType.INVALID_FIELD, 
validationError.getType());
+        assertEquals("/value", validationError.getFieldName().orElse(""));
+    }
+
+    @Test
+    public void testFieldValidatorReturningEmptyCollectionIsValid() {
+        final FieldValidator fieldValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                return List.of();
+            }
+
+            @Override
+            public String getDescription() {
+                return "always-valid validator";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"anything"));
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(FIELD_VALUE, List.of(fieldValidator)), List.of());
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, false, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertTrue(result.isValid());
+    }
+
+    @Test
+    public void testMultipleFieldValidatorsAllInvoked() {
+        final AtomicBoolean firstInvoked = new AtomicBoolean(false);
+        final AtomicBoolean secondInvoked = new AtomicBoolean(false);
+
+        final FieldValidator first = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                firstInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(path).explanation("first 
failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "first";
+            }
+        };
+
+        final FieldValidator second = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                secondInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(path).explanation("second 
failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "second";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"test"));
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(FIELD_VALUE, List.of(first, second)), List.of());
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, false, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertTrue(firstInvoked.get());
+        assertTrue(secondInvoked.get());
+        assertFalse(result.isValid());
+        assertEquals(2, result.getValidationErrors().size());
+    }
+
+    @Test
+    public void testFieldValidatorNotInvokedWhenValueNull() {
+        final AtomicBoolean invoked = new AtomicBoolean(false);
+        final FieldValidator fieldValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                invoked.set(true);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("should not be invoked")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "validator for 
testFieldValidatorNotInvokedWhenValueNull";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final Map<String, Object> values = new HashMap<>();
+        values.put(FIELD_VALUE, null);
+        final MapRecord record = new MapRecord(schema, values);
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(FIELD_VALUE, List.of(fieldValidator)), List.of());
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, false, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertFalse(invoked.get());
+        assertTrue(result.isValid());
+    }
+
+    @Test
+    public void testRecordValidatorIsApplied() {
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+
+        final RecordValidator recordValidator = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(fieldPath + "/" + FIELD_VALUE)
+                        .inputValue(record.getValue(FIELD_VALUE))
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("value must equal 'expected'")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "record validator for testRecordValidatorIsApplied";
+            }
+        };
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(), List.of(recordValidator));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"actual"));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, true, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertFalse(result.isValid());
+        final ValidationError validationError = 
result.getValidationErrors().iterator().next();
+        assertEquals("/value", validationError.getFieldName().orElse(""));
+        assertEquals(ValidationErrorType.INVALID_FIELD, 
validationError.getType());
+    }
+
+    @Test
+    public void testRecordValidatorReturningEmptyCollectionIsValid() {
+        final RecordValidator recordValidator = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                return List.of();
+            }
+
+            @Override
+            public String getDescription() {
+                return "always-valid record validator";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"anything"));
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(), List.of(recordValidator));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, true, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertTrue(result.isValid());
+    }
+
+    @Test
+    public void testMultipleRecordValidatorsAllInvoked() {
+        final AtomicBoolean firstInvoked = new AtomicBoolean(false);
+        final AtomicBoolean secondInvoked = new AtomicBoolean(false);
+
+        final RecordValidator first = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                firstInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(fieldPath).explanation("first
 record validator failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "first record validator";
+            }
+        };
+
+        final RecordValidator second = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                secondInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(fieldPath).explanation("second
 record validator failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "second record validator";
+            }
+        };
+
+        final RecordField recordField = new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType());
+        final RecordSchema schema = new 
SimpleRecordSchema(List.of(recordField));
+        final MapRecord record = new MapRecord(schema, Map.of(FIELD_VALUE, 
"test"));
+
+        final SchemaValidators schemaValidators = new 
SchemaValidators(Map.of(), List.of(first, second));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(schema, true, true, 
schemaValidators));
+        final SchemaValidationResult result = 
validatorService.validate(record);
+
+        assertTrue(firstInvoked.get());
+        assertTrue(secondInvoked.get());
+        assertFalse(result.isValid());
+        assertEquals(2, result.getValidationErrors().size());
+    }
+
+    @Test
+    public void testNestedRecordFieldValidatorIsApplied() {
+        final AtomicBoolean invoked = new AtomicBoolean(false);
+        final FieldValidator childNameValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                invoked.set(true);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("child name is invalid")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "child name validator";
+            }
+        };
+
+        final RecordSchema childSchema = new SimpleRecordSchema(List.of(new 
RecordField(FIELD_NAME, RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ID, RecordFieldType.INT.getDataType()),
+                new RecordField(FIELD_CHILD, 
RecordFieldType.RECORD.getRecordDataType(childSchema))));
+
+        final MapRecord childRecord = new MapRecord(childSchema, 
Map.of(FIELD_NAME, "test"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_ID, 1, FIELD_CHILD, childRecord));
+
+        final SchemaValidators childValidators = new 
SchemaValidators(Map.of(FIELD_NAME, List.of(childNameValidator)), List.of());
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(), Map.of(FIELD_CHILD, childValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertTrue(invoked.get());
+        assertFalse(result.isValid());
+        assertEquals(1, result.getValidationErrors().size());
+        final ValidationError error = 
result.getValidationErrors().iterator().next();
+        assertEquals("/child/name", error.getFieldName().orElse(""));
+    }
+
+    @Test
+    public void testArrayOfRecordFieldValidatorIsApplied() {
+        final List<String> capturedPaths = new ArrayList<>();
+        final FieldValidator keyValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                capturedPaths.add(path);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("key is invalid")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "key validator";
+            }
+        };
+
+        final RecordSchema elementSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_KEY, 
RecordFieldType.STRING.getDataType()),
+                new RecordField(FIELD_VALUE, 
RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ITEMS, 
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.RECORD.getRecordDataType(elementSchema)))));
+
+        final MapRecord element0 = new MapRecord(elementSchema, 
Map.of(FIELD_KEY, "k0", FIELD_VALUE, "v0"));
+        final MapRecord element1 = new MapRecord(elementSchema, 
Map.of(FIELD_KEY, "k1", FIELD_VALUE, "v1"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_ITEMS, new Object[]{element0, element1}));
+
+        final SchemaValidators elementValidators = new 
SchemaValidators(Map.of(FIELD_KEY, List.of(keyValidator)), List.of());
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(), Map.of(FIELD_ITEMS, elementValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertFalse(result.isValid());
+        assertEquals(2, result.getValidationErrors().size());
+        assertEquals(2, capturedPaths.size());
+        assertEquals("/items[0]/key", capturedPaths.get(0));
+        assertEquals("/items[1]/key", capturedPaths.get(1));
+    }
+
+    @Test
+    public void testFieldValidatorScopedToNestingLevel() {
+        final AtomicBoolean rootInvoked = new AtomicBoolean(false);
+        final AtomicBoolean childInvoked = new AtomicBoolean(false);
+
+        final FieldValidator rootNameValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                rootInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(path).explanation("root name 
failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "root name validator";
+            }
+        };
+
+        final FieldValidator childNameValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                childInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(path).explanation("child 
name failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "child name validator";
+            }
+        };
+
+        final RecordSchema childSchema = new SimpleRecordSchema(List.of(new 
RecordField(FIELD_NAME, RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_NAME, 
RecordFieldType.STRING.getDataType()),
+                new RecordField(FIELD_CHILD, 
RecordFieldType.RECORD.getRecordDataType(childSchema))));
+
+        final MapRecord childRecord = new MapRecord(childSchema, 
Map.of(FIELD_NAME, "childValue"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_NAME, "parentValue", FIELD_CHILD, childRecord));
+
+        final SchemaValidators childValidators = new 
SchemaValidators(Map.of(FIELD_NAME, List.of(childNameValidator)), List.of());
+        final SchemaValidators rootValidators = new 
SchemaValidators(Map.of(FIELD_NAME, List.of(rootNameValidator)), List.of(), 
Map.of(FIELD_CHILD, childValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertTrue(rootInvoked.get());
+        assertTrue(childInvoked.get());
+        assertFalse(result.isValid());
+        assertEquals(2, result.getValidationErrors().size());
+
+        final List<String> errorFields = result.getValidationErrors().stream()
+                .map(e -> e.getFieldName().orElse(""))
+                .sorted()
+                .collect(Collectors.toList());
+        assertEquals("/child/name", errorFields.get(0));
+        assertEquals("/name", errorFields.get(1));
+    }
+
+    @Test
+    public void testNestedRecordValidatorIsApplied() {
+        final AtomicBoolean rootRecordValidatorInvoked = new 
AtomicBoolean(false);
+        final AtomicBoolean childRecordValidatorInvoked = new 
AtomicBoolean(false);
+        final List<String> capturedChildPaths = new ArrayList<>();
+
+        final RecordValidator rootRecordValidator = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                rootRecordValidatorInvoked.set(true);
+                return List.of();
+            }
+
+            @Override
+            public String getDescription() {
+                return "root record validator";
+            }
+        };
+
+        final RecordValidator childRecordValidator = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                childRecordValidatorInvoked.set(true);
+                capturedChildPaths.add(fieldPath);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(fieldPath)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("child record is invalid")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "child record validator";
+            }
+        };
+
+        final RecordSchema childSchema = new SimpleRecordSchema(List.of(new 
RecordField(FIELD_VALUE, RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ID, RecordFieldType.INT.getDataType()),
+                new RecordField(FIELD_CHILD, 
RecordFieldType.RECORD.getRecordDataType(childSchema))));
+
+        final MapRecord childRecord = new MapRecord(childSchema, 
Map.of(FIELD_VALUE, "test"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_ID, 1, FIELD_CHILD, childRecord));
+
+        final SchemaValidators childValidators = new 
SchemaValidators(Map.of(), List.of(childRecordValidator));
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(rootRecordValidator), Map.of(FIELD_CHILD, childValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertTrue(rootRecordValidatorInvoked.get());
+        assertTrue(childRecordValidatorInvoked.get());
+        assertFalse(result.isValid());
+        assertEquals(1, result.getValidationErrors().size());
+        assertEquals(1, capturedChildPaths.size());
+        assertEquals("/child", capturedChildPaths.get(0));
+    }
+
+    @Test
+    public void testThreeLevelNestedValidators() {
+        final List<String> capturedPaths = new ArrayList<>();
+
+        final FieldValidator leafValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                capturedPaths.add(path);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("leaf value is invalid")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "leaf validator";
+            }
+        };
+
+        final String fieldGrandchild = "grandchild";
+        final String fieldCode = "code";
+
+        final RecordSchema grandchildSchema = new 
SimpleRecordSchema(List.of(new RecordField(fieldCode, 
RecordFieldType.STRING.getDataType())));
+        final RecordSchema childSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_NAME, 
RecordFieldType.STRING.getDataType()),
+                new RecordField(fieldGrandchild, 
RecordFieldType.RECORD.getRecordDataType(grandchildSchema))));
+        final RecordSchema rootSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ID, RecordFieldType.INT.getDataType()),
+                new RecordField(FIELD_CHILD, 
RecordFieldType.RECORD.getRecordDataType(childSchema))));
+
+        final MapRecord grandchildRecord = new MapRecord(grandchildSchema, 
Map.of(fieldCode, "ABC"));
+        final MapRecord childRecord = new MapRecord(childSchema, 
Map.of(FIELD_NAME, "test", fieldGrandchild, grandchildRecord));
+        final MapRecord rootRecord = new MapRecord(rootSchema, 
Map.of(FIELD_ID, 1, FIELD_CHILD, childRecord));
+
+        final SchemaValidators grandchildValidators = new 
SchemaValidators(Map.of(fieldCode, List.of(leafValidator)), List.of());
+        final SchemaValidators childValidators = new 
SchemaValidators(Map.of(), List.of(), Map.of(fieldGrandchild, 
grandchildValidators));
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(), Map.of(FIELD_CHILD, childValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(rootSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(rootRecord);
+
+        assertFalse(result.isValid());
+        assertEquals(1, result.getValidationErrors().size());
+        assertEquals(1, capturedPaths.size());
+        assertEquals("/child/grandchild/code", capturedPaths.get(0));
+    }
+
+    @Test
+    public void testArrayWithNullElementsSkipsValidation() {
+        final List<String> capturedPaths = new ArrayList<>();
+        final FieldValidator keyValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                capturedPaths.add(path);
+                return List.of(DefaultValidationError.builder()
+                        .fieldName(path)
+                        .inputValue(value)
+                        .type(ValidationErrorType.INVALID_FIELD)
+                        .explanation("key is invalid")
+                        .build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "key validator";
+            }
+        };
+
+        final RecordSchema elementSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_KEY, 
RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ITEMS, 
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.RECORD.getRecordDataType(elementSchema),
 true))));
+
+        final MapRecord element0 = new MapRecord(elementSchema, 
Map.of(FIELD_KEY, "k0"));
+        final MapRecord element2 = new MapRecord(elementSchema, 
Map.of(FIELD_KEY, "k2"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_ITEMS, new Object[]{element0, null, element2}));
+
+        final SchemaValidators elementValidators = new 
SchemaValidators(Map.of(FIELD_KEY, List.of(keyValidator)), List.of());
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(), Map.of(FIELD_ITEMS, elementValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertFalse(result.isValid());
+        assertEquals(2, capturedPaths.size());
+        assertEquals("/items[0]/key", capturedPaths.get(0));
+        assertEquals("/items[2]/key", capturedPaths.get(1));
+    }
+
+    @Test
+    public void testCombinedFieldAndRecordValidatorsOnNestedRecord() {
+        final AtomicBoolean fieldValidatorInvoked = new AtomicBoolean(false);
+        final AtomicBoolean recordValidatorInvoked = new AtomicBoolean(false);
+
+        final FieldValidator fieldValidator = new FieldValidator() {
+            @Override
+            public Collection<ValidationError> validate(final String path, 
final Object value) {
+                fieldValidatorInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(path).explanation("field 
failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "child field validator";
+            }
+        };
+
+        final RecordValidator recordValidator = new RecordValidator() {
+            @Override
+            public Collection<ValidationError> validate(final Record record, 
final String fieldPath) {
+                recordValidatorInvoked.set(true);
+                return 
List.of(DefaultValidationError.builder().fieldName(fieldPath).explanation("record
 failed").build());
+            }
+
+            @Override
+            public String getDescription() {
+                return "child record validator";
+            }
+        };
+
+        final RecordSchema childSchema = new SimpleRecordSchema(List.of(new 
RecordField(FIELD_NAME, RecordFieldType.STRING.getDataType())));
+        final RecordSchema parentSchema = new SimpleRecordSchema(List.of(
+                new RecordField(FIELD_ID, RecordFieldType.INT.getDataType()),
+                new RecordField(FIELD_CHILD, 
RecordFieldType.RECORD.getRecordDataType(childSchema))));
+
+        final MapRecord childRecord = new MapRecord(childSchema, 
Map.of(FIELD_NAME, "test"));
+        final MapRecord parentRecord = new MapRecord(parentSchema, 
Map.of(FIELD_ID, 1, FIELD_CHILD, childRecord));
+
+        final SchemaValidators childValidators = new 
SchemaValidators(Map.of(FIELD_NAME, List.of(fieldValidator)), 
List.of(recordValidator));
+        final SchemaValidators rootValidators = new SchemaValidators(Map.of(), 
List.of(), Map.of(FIELD_CHILD, childValidators));
+        final StandardSchemaValidator validatorService = new 
StandardSchemaValidator(new SchemaValidationContext(parentSchema, true, true, 
rootValidators));
+        final SchemaValidationResult result = 
validatorService.validate(parentRecord);
+
+        assertTrue(fieldValidatorInvoked.get());
+        assertTrue(recordValidatorInvoked.get());
+        assertFalse(result.isValid());
+        assertEquals(2, result.getValidationErrors().size());
+    }
+
     private void whenValueIsAcceptedAsDataTypeThenConsideredAsValid(final 
Object value, final RecordFieldType schemaDataType) {
         final SchemaValidationResult result = whenSingleValueIsTested(value, 
schemaDataType);
         thenSingleValueIsValid(result);

Reply via email to