This is an automated email from the ASF dual-hosted git repository.

pvillard31 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 1a22684df8d NIFI-16187 FlowFile grouping strategy in ConsumeKafka 
(#11537)
1a22684df8d is described below

commit 1a22684df8d7ea3aa8bb0952c001894b16c85d0a
Author: Alaksiej Ščarbaty <[email protected]>
AuthorDate: Thu Aug 20 12:32:19 2026 +0200

    NIFI-16187 FlowFile grouping strategy in ConsumeKafka (#11537)
    
    Extract FlowFile grouping into RecordGroupingStrategy, then add 
MergeSchemaGrouping
    and a Schema Conflict Resolution property so RECORD consumption can batch 
different
    schemas into one FlowFile per topic/partition group while preserving Create 
New
    FlowFile as the default.
    
    Make DataTypeUtils.merge nullify fields present in only one schema so a 
merged
    schema can describe records from either input without failing strict 
writers.
---
 .../serialization/record/util/DataTypeUtils.java   |  33 ++-
 .../serialization/record/TestDataTypeUtils.java    | 308 ++++++++++++++++++++-
 .../kafka/processors/AbstractConsumeKafkaIT.java   |  17 ++
 .../processors/ConsumeKafkaMergeSchemaIT.java      | 301 ++++++++++++++++++++
 .../processors/ConsumeKafkaMergeSchemaTypedIT.java | 215 ++++++++++++++
 .../apache/nifi/kafka/processors/ConsumeKafka.java |  39 ++-
 .../AbstractRecordStreamKafkaMessageConverter.java | 148 +---------
 .../convert/CreateNewFlowFileGrouping.java         | 166 +++++++++++
 ...ectOffsetRecordStreamKafkaMessageConverter.java |   6 +-
 .../consumer/convert/MergeSchemaGrouping.java      | 154 +++++++++++
 .../consumer/convert/RecordGroupingStrategy.java   |  47 ++++
 .../convert/RecordStreamKafkaMessageConverter.java |   6 +-
 .../WrapperRecordStreamKafkaMessageConverter.java  |   6 +-
 .../convert/CreateNewFlowFileGroupingTest.java     | 165 +++++++++++
 .../consumer/convert/MergeSchemaGroupingTest.java  | 260 +++++++++++++++++
 .../RecordStreamKafkaMessageConverterTest.java     |   3 +-
 ...apperRecordStreamKafkaMessageConverterTest.java |   3 +-
 .../shared/property/SchemaConflictResolution.java  |  55 ++++
 18 files changed, 1770 insertions(+), 162 deletions(-)

diff --git 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/util/DataTypeUtils.java
 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/util/DataTypeUtils.java
index a8c1f194f45..37751ef125a 100644
--- 
a/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/util/DataTypeUtils.java
+++ 
b/nifi-commons/nifi-record/src/main/java/org/apache/nifi/serialization/record/util/DataTypeUtils.java
@@ -59,6 +59,7 @@ import java.time.format.DateTimeFormatter;
 import java.time.format.DateTimeParseException;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.BitSet;
 import java.util.Collection;
 import java.util.Comparator;
 import java.util.EnumMap;
@@ -1603,14 +1604,7 @@ public class DataTypeUtils {
         }
 
         final List<RecordField> otherFields = otherSchema.getFields();
-        if (otherFields.isEmpty()) {
-            return thisSchema;
-        }
-
         final List<RecordField> thisFields = thisSchema.getFields();
-        if (thisFields.isEmpty()) {
-            return otherSchema;
-        }
 
         final Map<String, Integer> fieldIndices = new HashMap<>();
         final List<RecordField> fields = new ArrayList<>();
@@ -1627,6 +1621,7 @@ public class DataTypeUtils {
             fields.add(field);
         }
 
+        final BitSet matchedFieldIndices = new BitSet(thisFields.size());
         for (final RecordField otherField : otherFields) {
             Integer fieldIndex = fieldIndices.get(otherField.getFieldName());
 
@@ -1642,11 +1637,15 @@ public class DataTypeUtils {
             }
 
             // If there is no field with the same name then just add 
'otherField'.
+            // Fields present in only one schema are nullable in the merged 
schema,
+            // since the merged schema is a superset of both inputs.
             if (fieldIndex == null) {
-                fields.add(otherField);
+                fields.add(makeNullable(otherField));
                 continue;
             }
 
+            matchedFieldIndices.set(fieldIndex);
+
             // Merge the two fields, if necessary
             final RecordField thisField = fields.get(fieldIndex);
             if (isMergeRequired(thisField, otherField)) {
@@ -1655,6 +1654,12 @@ public class DataTypeUtils {
             }
         }
 
+        for (int i = 0; i < thisFields.size(); i++) {
+            if (!matchedFieldIndices.get(i)) {
+                fields.set(i, makeNullable(fields.get(i)));
+            }
+        }
+
         return new SimpleRecordSchema(fields);
     }
 
@@ -1667,9 +1672,21 @@ public class DataTypeUtils {
             return true;
         }
 
+        if (thisField.isNullable() != otherField.isNullable()) {
+            return true;
+        }
+
         return !Objects.equals(thisField.getDefaultValue(), 
otherField.getDefaultValue());
     }
 
+    private static RecordField makeNullable(final RecordField field) {
+        if (field.isNullable()) {
+            return field;
+        }
+
+        return new RecordField(field.getFieldName(), field.getDataType(), 
field.getDefaultValue(), field.getAliases(), true);
+    }
+
     public static RecordField merge(final RecordField thisField, final 
RecordField otherField) {
         final String fieldName = thisField.getFieldName();
         final Set<String> aliases = new HashSet<>();
diff --git 
a/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/TestDataTypeUtils.java
 
b/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/TestDataTypeUtils.java
index ab5624e392a..ac358c2c800 100644
--- 
a/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/TestDataTypeUtils.java
+++ 
b/nifi-commons/nifi-record/src/test/java/org/apache/nifi/serialization/record/TestDataTypeUtils.java
@@ -47,6 +47,7 @@ import java.util.LinkedList;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
+import java.util.Set;
 import java.util.UUID;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.DoubleAdder;
@@ -1298,14 +1299,14 @@ public class TestDataTypeUtils {
     @Test
     public void testMergeDataTypesMergesRecordSchemasInsteadOfCreatingChoice() 
{
         final RecordSchema schemaA = new SimpleRecordSchema(List.of(
-            new RecordField("firstName", RecordFieldType.STRING.getDataType()),
-            new RecordField("lastName", RecordFieldType.STRING.getDataType()),
-            new RecordField("address", RecordFieldType.STRING.getDataType())));
+            new RecordField("firstName", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("lastName", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("address", RecordFieldType.STRING.getDataType(), 
false)));
 
         final RecordSchema schemaB = new SimpleRecordSchema(List.of(
-            new RecordField("firstName", RecordFieldType.STRING.getDataType()),
-            new RecordField("lastName", RecordFieldType.STRING.getDataType()),
-            new RecordField("age", RecordFieldType.INT.getDataType())));
+            new RecordField("firstName", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("lastName", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("age", RecordFieldType.INT.getDataType(), false)));
 
         final DataType recordTypeA = 
RecordFieldType.RECORD.getRecordDataType(schemaA);
         final DataType recordTypeB = 
RecordFieldType.RECORD.getRecordDataType(schemaB);
@@ -1319,6 +1320,10 @@ public class TestDataTypeUtils {
         assertTrue(mergedSchema.getField("lastName").isPresent());
         assertTrue(mergedSchema.getField("address").isPresent());
         assertTrue(mergedSchema.getField("age").isPresent());
+        
assertFalse(mergedSchema.getField("firstName").orElseThrow().isNullable());
+        
assertFalse(mergedSchema.getField("lastName").orElseThrow().isNullable());
+        
assertTrue(mergedSchema.getField("address").orElseThrow().isNullable());
+        assertTrue(mergedSchema.getField("age").orElseThrow().isNullable());
     }
 
     @Test
@@ -1392,4 +1397,295 @@ public class TestDataTypeUtils {
         assertTrue(finalSchema.getField("field_0").isPresent());
         assertTrue(finalSchema.getField("field_4999").isPresent());
     }
+
+    @Test
+    public void testMergeSchemasMakesSingleSideFieldsNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("name", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(2, merged.getFieldCount());
+        assertTrue(merged.getField("id").orElseThrow().isNullable());
+        assertTrue(merged.getField("name").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasWhenOneSchemaIsStrictSubset() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("x", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("shared", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("shared", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(2, merged.getFieldCount());
+        assertTrue(merged.getField("x").orElseThrow().isNullable());
+        assertFalse(merged.getField("shared").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasKeepsSharedNonNullableFieldNonNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(1, merged.getFieldCount());
+        assertFalse(merged.getField("id").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasPreservesNullableWhenLeftIsNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
true)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertTrue(merged.getField("id").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasPreservesNullableWhenRightIsNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
true)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertTrue(merged.getField("id").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasWidensTypesAndNullifiesSingleSideFields() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.INT.getDataType(), false),
+            new RecordField("onlyA", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.LONG.getDataType(), true),
+            new RecordField("onlyB", RecordFieldType.STRING.getDataType(), 
Set.of("bAlias"), false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(3, merged.getFieldCount());
+        final RecordField idField = merged.getField("id").orElseThrow();
+        assertEquals(RecordFieldType.LONG, 
idField.getDataType().getFieldType());
+        assertTrue(idField.isNullable());
+        assertTrue(merged.getField("onlyA").orElseThrow().isNullable());
+        assertTrue(merged.getField("onlyB").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasIdenticalSchemasRetainNullability() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("required", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("optional", RecordFieldType.INT.getDataType(), 
true)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("required", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("optional", RecordFieldType.INT.getDataType(), 
true)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(2, merged.getFieldCount());
+        assertFalse(merged.getField("required").orElseThrow().isNullable());
+        assertTrue(merged.getField("optional").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasMakesNestedSingleSideFieldsNullable() {
+        final RecordSchema nestedA = new SimpleRecordSchema(List.of(
+            new RecordField("street", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("city", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema nestedB = new SimpleRecordSchema(List.of(
+            new RecordField("street", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("zip", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("address", 
RecordFieldType.RECORD.getRecordDataType(nestedA), false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("address", 
RecordFieldType.RECORD.getRecordDataType(nestedB), false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+        final RecordField addressField = 
merged.getField("address").orElseThrow();
+        assertFalse(addressField.isNullable());
+
+        final RecordSchema nestedMerged = ((RecordDataType) 
addressField.getDataType()).getChildSchema();
+        assertEquals(3, nestedMerged.getFieldCount());
+        
assertFalse(nestedMerged.getField("street").orElseThrow().isNullable());
+        assertTrue(nestedMerged.getField("city").orElseThrow().isNullable());
+        assertTrue(nestedMerged.getField("zip").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasMakesArrayOfRecordSingleSideFieldsNullable() {
+        final RecordSchema elementA = new SimpleRecordSchema(List.of(
+            new RecordField("x", RecordFieldType.INT.getDataType(), false),
+            new RecordField("shared", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema elementB = new SimpleRecordSchema(List.of(
+            new RecordField("y", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("shared", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("items", 
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.RECORD.getRecordDataType(elementA)),
 false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("items", 
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.RECORD.getRecordDataType(elementB)),
 false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+        final RecordField itemsField = merged.getField("items").orElseThrow();
+        assertFalse(itemsField.isNullable());
+
+        final DataType elementType = ((ArrayDataType) 
itemsField.getDataType()).getElementType();
+        final RecordSchema elementMerged = ((RecordDataType) 
elementType).getChildSchema();
+        assertEquals(3, elementMerged.getFieldCount());
+        
assertFalse(elementMerged.getField("shared").orElseThrow().isNullable());
+        assertTrue(elementMerged.getField("x").orElseThrow().isNullable());
+        assertTrue(elementMerged.getField("y").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void 
testMergeSchemasNullifiesTopLevelRecordFieldPresentOnOnlyOneSide() {
+        final RecordSchema nested = new SimpleRecordSchema(List.of(
+            new RecordField("value", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("details", 
RecordFieldType.RECORD.getRecordDataType(nested), false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertFalse(merged.getField("id").orElseThrow().isNullable());
+        assertTrue(merged.getField("details").orElseThrow().isNullable());
+        final RecordSchema detailsSchema = ((RecordDataType) 
merged.getField("details").orElseThrow().getDataType()).getChildSchema();
+        
assertFalse(detailsSchema.getField("value").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasMakesDeeplyNestedSingleSideFieldsNullable() {
+        final RecordSchema leafA = new SimpleRecordSchema(List.of(
+            new RecordField("aOnly", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("common", RecordFieldType.INT.getDataType(), 
false)));
+        final RecordSchema leafB = new SimpleRecordSchema(List.of(
+            new RecordField("bOnly", RecordFieldType.STRING.getDataType(), 
false),
+            new RecordField("common", RecordFieldType.INT.getDataType(), 
false)));
+
+        final RecordSchema midA = new SimpleRecordSchema(List.of(
+            new RecordField("leaf", 
RecordFieldType.RECORD.getRecordDataType(leafA), false)));
+        final RecordSchema midB = new SimpleRecordSchema(List.of(
+            new RecordField("leaf", 
RecordFieldType.RECORD.getRecordDataType(leafB), false)));
+
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("mid", 
RecordFieldType.RECORD.getRecordDataType(midA), false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("mid", 
RecordFieldType.RECORD.getRecordDataType(midB), false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+        final RecordSchema midMerged = ((RecordDataType) 
merged.getField("mid").orElseThrow().getDataType()).getChildSchema();
+        final RecordSchema leafMerged = ((RecordDataType) 
midMerged.getField("leaf").orElseThrow().getDataType()).getChildSchema();
+
+        assertFalse(leafMerged.getField("common").orElseThrow().isNullable());
+        assertTrue(leafMerged.getField("aOnly").orElseThrow().isNullable());
+        assertTrue(leafMerged.getField("bOnly").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasMatchesFieldsByAliasWithoutNullifying() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
Set.of("identifier"), false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(1, merged.getFieldCount());
+        final RecordField mergedField = merged.getField("id").orElseThrow();
+        assertFalse(mergedField.isNullable());
+        assertTrue(mergedField.getAliases().contains("identifier"));
+    }
+
+    @Test
+    public void testMergeSchemasWithEmptySchemaMakesOtherFieldsNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema empty = new SimpleRecordSchema(List.of());
+
+        final RecordSchema mergedWithEmptyOther = DataTypeUtils.merge(schemaA, 
empty);
+        final RecordSchema mergedWithEmptyThis = DataTypeUtils.merge(empty, 
schemaA);
+
+        assertEquals(1, mergedWithEmptyOther.getFieldCount());
+        
assertTrue(mergedWithEmptyOther.getField("id").orElseThrow().isNullable());
+        assertEquals(1, mergedWithEmptyThis.getFieldCount());
+        
assertTrue(mergedWithEmptyThis.getField("id").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void 
testMergeSchemasPreservesDefaultValueWhenNullifyingSingleSideField() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("onlyA", RecordFieldType.STRING.getDataType(), 
"keep-me", false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("onlyB", RecordFieldType.STRING.getDataType(), 
"also-keep", false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertTrue(merged.getField("onlyA").orElseThrow().isNullable());
+        assertEquals("keep-me", 
merged.getField("onlyA").orElseThrow().getDefaultValue());
+        assertTrue(merged.getField("onlyB").orElseThrow().isNullable());
+        assertEquals("also-keep", 
merged.getField("onlyB").orElseThrow().getDefaultValue());
+    }
+
+    @Test
+    public void testMergeSchemasMatchesFieldUsingLeftAlias() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("identifier", 
RecordFieldType.STRING.getDataType(), Set.of("id"), false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(1, merged.getFieldCount());
+        final RecordField mergedField = 
merged.getField("identifier").orElseThrow();
+        assertFalse(mergedField.isNullable());
+        assertTrue(mergedField.getAliases().contains("id"));
+    }
+
+    @Test
+    public void 
testMergeSchemasPreservesLeftFieldOrderThenAppendsRightOnlyFields() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("b", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("a", RecordFieldType.STRING.getDataType(), 
false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("c", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("a", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertEquals(List.of("b", "a", "c"), merged.getFieldNames());
+        assertTrue(merged.getField("b").orElseThrow().isNullable());
+        assertFalse(merged.getField("a").orElseThrow().isNullable());
+        assertTrue(merged.getField("c").orElseThrow().isNullable());
+    }
+
+    @Test
+    public void testMergeSchemasSingleSideNestedArrayFieldBecomesNullable() {
+        final RecordSchema schemaA = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), false),
+            new RecordField("tags", 
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.STRING.getDataType(), 
false), false)));
+        final RecordSchema schemaB = new SimpleRecordSchema(List.of(
+            new RecordField("id", RecordFieldType.STRING.getDataType(), 
false)));
+
+        final RecordSchema merged = DataTypeUtils.merge(schemaA, schemaB);
+
+        assertFalse(merged.getField("id").orElseThrow().isNullable());
+        assertTrue(merged.getField("tags").orElseThrow().isNullable());
+        assertFalse(((ArrayDataType) 
merged.getField("tags").orElseThrow().getDataType()).isElementsNullable());
+    }
 }
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/AbstractConsumeKafkaIT.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/AbstractConsumeKafkaIT.java
index 577cc0a6246..c4e481b7822 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/AbstractConsumeKafkaIT.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/AbstractConsumeKafkaIT.java
@@ -22,17 +22,21 @@ import org.apache.kafka.clients.producer.ProducerRecord;
 import org.apache.kafka.clients.producer.RecordMetadata;
 import org.apache.kafka.common.header.Header;
 import org.apache.kafka.common.serialization.StringSerializer;
+import org.apache.nifi.util.TestRunner;
 
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.List;
 import java.util.Properties;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.Future;
+import java.util.function.Predicate;
 import java.util.stream.Collectors;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.fail;
 
 public abstract class AbstractConsumeKafkaIT extends AbstractKafkaBaseIT {
 
@@ -72,4 +76,17 @@ public abstract class AbstractConsumeKafkaIT extends 
AbstractKafkaBaseIT {
             assertTrue(metadata.hasOffset());
         }
     }
+
+    /**
+     * Runs the processor until {@code condition} is satisfied or {@code 
timeout} elapses.
+     */
+    protected void runUntil(final TestRunner runner, final 
Predicate<TestRunner> condition, final Duration timeout) {
+        final long deadline = System.nanoTime() + timeout.toNanos();
+        while (!condition.test(runner)) {
+            if (System.nanoTime() >= deadline) {
+                fail("Timed out after " + timeout + " waiting for condition to 
be met");
+            }
+            runner.run(1, false, false);
+        }
+    }
 }
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaIT.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaIT.java
new file mode 100644
index 00000000000..ba6770ba0a4
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaIT.java
@@ -0,0 +1,301 @@
+/*
+ * 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.kafka.processors;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.node.ArrayNode;
+import org.apache.kafka.clients.admin.AdminClient;
+import org.apache.kafka.clients.admin.AdminClientConfig;
+import org.apache.kafka.clients.admin.NewTopic;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.header.Header;
+import org.apache.nifi.kafka.processors.consumer.ProcessingStrategy;
+import org.apache.nifi.kafka.service.api.consumer.AutoOffsetReset;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.kafka.shared.property.OutputStrategy;
+import org.apache.nifi.kafka.shared.property.SchemaConflictResolution;
+import org.apache.nifi.reporting.InitializationException;
+import org.apache.nifi.util.MockFlowFile;
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.time.Duration;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.ExecutionException;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class ConsumeKafkaMergeSchemaIT extends AbstractConsumeKafkaIT {
+
+    private static final int FIRST_PARTITION = 0;
+
+    private static final String RECORD_WITH_ID = """
+            { "id": 1 }
+            """;
+
+    private static final String RECORD_WITH_NAME = """
+            { "name": "Alice" }
+            """;
+
+    private static final String INVALID_RECORD = "not-valid-json";
+
+    private TestRunner runner;
+
+    @BeforeEach
+    void setRunner() throws InitializationException {
+        runner = TestRunners.newTestRunner(ConsumeKafka.class);
+        addKafkaConnectionService(runner);
+        runner.setProperty(ConsumeKafka.CONNECTION_SERVICE, 
CONNECTION_SERVICE_ID);
+        addRecordReaderService(runner);
+        addRecordWriterService(runner);
+    }
+
+    @Test
+    void testMergedSchemaProducesSingleFlowFile() throws ExecutionException, 
InterruptedException, IOException {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
OutputStrategy.USE_VALUE.getValue());
+        runner.setProperty(ConsumeKafka.SCHEMA_CONFLICT_RESOLUTION, 
SchemaConflictResolution.CONTINUE_WITH_MERGED_SCHEMA.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        produce(topic, List.of(
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_ID, List.<Header>of()),
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_NAME, List.<Header>of())));
+
+        runUntil(runner, r -> 
!r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS).isEmpty(), 
Duration.ofSeconds(30));
+
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, successFlowFiles.size());
+
+        final MockFlowFile flowFile = successFlowFiles.getFirst();
+        flowFile.assertAttributeEquals(KafkaFlowFileAttribute.KAFKA_TOPIC, 
topic);
+        flowFile.assertAttributeEquals(KafkaFlowFileAttribute.KAFKA_PARTITION, 
Integer.toString(FIRST_PARTITION));
+        flowFile.assertAttributeEquals("record.count", "2");
+
+        final JsonNode jsonTree = objectMapper.readTree(flowFile.getContent());
+        final JsonNode expected = objectMapper.readTree("""
+                [
+                  { "id": 1, "name": null },
+                  { "id": null, "name": "Alice" }
+                ]
+                """);
+        assertEquals(expected, jsonTree);
+    }
+
+    @Test
+    void testCreateNewFlowFileDefaultProducesMultipleFlowFiles() throws 
ExecutionException, InterruptedException {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
OutputStrategy.USE_VALUE.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        produce(topic, List.of(
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_ID, List.<Header>of()),
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_NAME, List.<Header>of())));
+
+        runUntil(runner, r -> 
r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS).size() >= 2, 
Duration.ofSeconds(30));
+
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, successFlowFiles.size());
+
+        for (final MockFlowFile flowFile : successFlowFiles) {
+            flowFile.assertAttributeEquals(KafkaFlowFileAttribute.KAFKA_TOPIC, 
topic);
+            
flowFile.assertAttributeEquals(KafkaFlowFileAttribute.KAFKA_PARTITION, 
Integer.toString(FIRST_PARTITION));
+        }
+    }
+
+    @Test
+    void testMergedSchemaDifferentPartitionsProduceSeparateFlowFiles() throws 
Exception {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        try (final AdminClient admin = AdminClient.create(
+                Map.of(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, 
kafkaContainer.getBootstrapServers()))) {
+            admin.createTopics(List.of(new NewTopic(topic, 2, (short) 
1))).all().get();
+        }
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
OutputStrategy.USE_VALUE.getValue());
+        runner.setProperty(ConsumeKafka.SCHEMA_CONFLICT_RESOLUTION, 
SchemaConflictResolution.CONTINUE_WITH_MERGED_SCHEMA.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        // Publish as one producer batch so the records are available together 
for a single poll/onTrigger.
+        produce(topic, List.of(
+                new ProducerRecord<>(topic, 0, (String) null, RECORD_WITH_ID, 
List.<Header>of()),
+                new ProducerRecord<>(topic, 0, (String) null, 
RECORD_WITH_NAME, List.<Header>of()),
+                new ProducerRecord<>(topic, 1, (String) null, 
RECORD_WITH_NAME, List.<Header>of())));
+
+        runUntil(runner, r -> 
totalRecordCount(r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS)) >= 3, 
Duration.ofSeconds(30));
+
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(3, totalRecordCount(successFlowFiles));
+
+        final List<MockFlowFile> partitionZero = successFlowFiles.stream()
+                .filter(ff -> 
"0".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION)))
+                .toList();
+        final List<MockFlowFile> partitionOne = successFlowFiles.stream()
+                .filter(ff -> 
"1".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION)))
+                .toList();
+
+        assertEquals(1, partitionZero.size());
+        assertEquals(1, partitionOne.size());
+        assertEquals(2, totalRecordCount(partitionZero));
+        assertEquals(1, totalRecordCount(partitionOne));
+        // Records from different partitions never share a FlowFile.
+        assertTrue(partitionZero.stream().noneMatch(ff -> 
"1".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION))));
+    }
+
+    private static int totalRecordCount(final List<MockFlowFile> flowFiles) {
+        return flowFiles.stream()
+                .mapToInt(ff -> 
Integer.parseInt(ff.getAttribute("record.count")))
+                .sum();
+    }
+
+    @Test
+    void testMergedSchemaWithParseFailure() throws ExecutionException, 
InterruptedException, IOException {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
OutputStrategy.USE_VALUE.getValue());
+        runner.setProperty(ConsumeKafka.SCHEMA_CONFLICT_RESOLUTION, 
SchemaConflictResolution.CONTINUE_WITH_MERGED_SCHEMA.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        produce(topic, List.of(
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_ID, List.<Header>of()),
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
INVALID_RECORD, List.<Header>of()),
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_NAME, List.<Header>of())));
+
+        runUntil(runner, r -> 
!r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS).isEmpty(), 
Duration.ofSeconds(30));
+
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, successFlowFiles.size());
+
+        final MockFlowFile successFlowFile = successFlowFiles.getFirst();
+        final JsonNode jsonTree = 
objectMapper.readTree(successFlowFile.getContent());
+        assertInstanceOf(ArrayNode.class, jsonTree);
+        assertEquals(2, jsonTree.size());
+
+        final List<MockFlowFile> parseFailureFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.PARSE_FAILURE);
+        assertEquals(1, parseFailureFlowFiles.size());
+        parseFailureFlowFiles.getFirst().assertContentEquals(INVALID_RECORD);
+    }
+
+    @Test
+    void testMergedSchemaWithInjectOffset() throws Exception {
+        final MockFlowFile flowFile = 
runMergedSchemaWithOutputStrategy(OutputStrategy.INJECT_OFFSET);
+        flowFile.assertAttributeEquals("record.count", "2");
+
+        final JsonNode jsonTree = objectMapper.readTree(flowFile.getContent());
+        assertInstanceOf(ArrayNode.class, jsonTree);
+        assertEquals(2, jsonTree.size());
+        assertEquals(1, jsonTree.get(0).get("id").asInt());
+        assertTrue(jsonTree.get(0).has("kafkaOffset"));
+        assertEquals("Alice", jsonTree.get(1).get("name").asText());
+        assertTrue(jsonTree.get(1).has("kafkaOffset"));
+    }
+
+    @Test
+    void testMergedSchemaWithUseWrapper() throws Exception {
+        final MockFlowFile flowFile = 
runMergedSchemaWithOutputStrategy(OutputStrategy.USE_WRAPPER);
+        flowFile.assertAttributeEquals("record.count", "2");
+
+        final JsonNode jsonTree = objectMapper.readTree(flowFile.getContent());
+        assertInstanceOf(ArrayNode.class, jsonTree);
+        assertEquals(2, jsonTree.size());
+        assertEquals(1, jsonTree.get(0).get("value").get("id").asInt());
+        assertTrue(jsonTree.get(0).has("metadata"));
+        assertEquals("Alice", 
jsonTree.get(1).get("value").get("name").asText());
+        assertTrue(jsonTree.get(1).has("metadata"));
+    }
+
+    @Test
+    void testMergedSchemaWithInjectMetadata() throws Exception {
+        final MockFlowFile flowFile = 
runMergedSchemaWithOutputStrategy(OutputStrategy.INJECT_METADATA);
+        flowFile.assertAttributeEquals("record.count", "2");
+
+        final JsonNode jsonTree = objectMapper.readTree(flowFile.getContent());
+        assertInstanceOf(ArrayNode.class, jsonTree);
+        assertEquals(2, jsonTree.size());
+        assertEquals(1, jsonTree.get(0).get("id").asInt());
+        assertTrue(jsonTree.get(0).has("kafkaMetadata"));
+        assertEquals("Alice", jsonTree.get(1).get("name").asText());
+        assertTrue(jsonTree.get(1).has("kafkaMetadata"));
+    }
+
+    private MockFlowFile runMergedSchemaWithOutputStrategy(final 
OutputStrategy outputStrategy)
+            throws ExecutionException, InterruptedException {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
outputStrategy.getValue());
+        runner.setProperty(ConsumeKafka.SCHEMA_CONFLICT_RESOLUTION, 
SchemaConflictResolution.CONTINUE_WITH_MERGED_SCHEMA.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        produce(topic, List.of(
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_ID, List.<Header>of()),
+                new ProducerRecord<>(topic, FIRST_PARTITION, (String) null, 
RECORD_WITH_NAME, List.<Header>of())));
+
+        runUntil(runner, r -> 
!r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS).isEmpty(), 
Duration.ofSeconds(30));
+
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, successFlowFiles.size());
+        return successFlowFiles.getFirst();
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaTypedIT.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaTypedIT.java
new file mode 100644
index 00000000000..8c7ef09ebab
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-3-integration/src/test/java/org/apache/nifi/kafka/processors/ConsumeKafkaMergeSchemaTypedIT.java
@@ -0,0 +1,215 @@
+/*
+ * 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.kafka.processors;
+
+import org.apache.avro.Schema;
+import org.apache.avro.file.DataFileStream;
+import org.apache.avro.file.DataFileWriter;
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericDatumReader;
+import org.apache.avro.generic.GenericDatumWriter;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.clients.producer.RecordMetadata;
+import org.apache.kafka.common.serialization.ByteArraySerializer;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.apache.nifi.avro.AvroReader;
+import org.apache.nifi.avro.AvroRecordSetWriter;
+import org.apache.nifi.avro.AvroTypeUtil;
+import org.apache.nifi.kafka.processors.consumer.ProcessingStrategy;
+import org.apache.nifi.kafka.service.api.consumer.AutoOffsetReset;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.kafka.shared.property.OutputStrategy;
+import org.apache.nifi.kafka.shared.property.SchemaConflictResolution;
+import org.apache.nifi.reporting.InitializationException;
+import org.apache.nifi.schema.access.SchemaAccessUtils;
+import org.apache.nifi.serialization.RecordReaderFactory;
+import org.apache.nifi.serialization.RecordSetWriterFactory;
+import org.apache.nifi.util.MockFlowFile;
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.UUID;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Future;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Verifies Continue with Merged Schema using typed Avro schemas (no 
inference):
+ * disjoint non-nullable fields become nullable so both records can be written 
together.
+ */
+class ConsumeKafkaMergeSchemaTypedIT extends AbstractConsumeKafkaIT {
+
+    private static final int FIRST_PARTITION = 0;
+
+    private static final String SCHEMA_WITH_ID = """
+            {
+              "type": "record",
+              "name": "IdRecord",
+              "fields": [
+                { "name": "id", "type": "long" }
+              ]
+            }
+            """;
+
+    private static final String SCHEMA_WITH_NAME = """
+            {
+              "type": "record",
+              "name": "NameRecord",
+              "fields": [
+                { "name": "name", "type": "string" }
+              ]
+            }
+            """;
+
+    private TestRunner runner;
+
+    @BeforeEach
+    void setRunner() throws InitializationException {
+        runner = TestRunners.newTestRunner(ConsumeKafka.class);
+        addKafkaConnectionService(runner);
+        runner.setProperty(ConsumeKafka.CONNECTION_SERVICE, 
CONNECTION_SERVICE_ID);
+        addEmbeddedAvroReader(runner);
+        addInheritAvroWriter(runner);
+    }
+
+    private void addEmbeddedAvroReader(final TestRunner runner) throws 
InitializationException {
+        final String readerId = ConsumeKafka.RECORD_READER.getName();
+        final RecordReaderFactory readerService = new AvroReader();
+        runner.addControllerService(readerId, readerService);
+        runner.setProperty(readerService, 
SchemaAccessUtils.SCHEMA_ACCESS_STRATEGY, "embedded-avro-schema");
+        runner.enableControllerService(readerService);
+        runner.setProperty(readerId, readerId);
+    }
+
+    private void addInheritAvroWriter(final TestRunner runner) throws 
InitializationException {
+        final String writerId = ConsumeKafka.RECORD_WRITER.getName();
+        final RecordSetWriterFactory writerService = new AvroRecordSetWriter();
+        runner.addControllerService(writerId, writerService);
+        runner.setProperty(writerService, 
SchemaAccessUtils.SCHEMA_ACCESS_STRATEGY, 
SchemaAccessUtils.INHERIT_RECORD_SCHEMA.getValue());
+        runner.enableControllerService(writerService);
+        runner.setProperty(writerId, writerId);
+    }
+
+    @Test
+    void testMergedSchemaWithDisjointNonNullableAvroFields() throws Exception {
+        final String topic = UUID.randomUUID().toString();
+        final String groupId = topic.substring(0, topic.indexOf("-"));
+
+        runner.setProperty(ConsumeKafka.GROUP_ID, groupId);
+        runner.setProperty(ConsumeKafka.TOPICS, topic);
+        runner.setProperty(ConsumeKafka.PROCESSING_STRATEGY, 
ProcessingStrategy.RECORD.getValue());
+        runner.setProperty(ConsumeKafka.OUTPUT_STRATEGY, 
OutputStrategy.USE_VALUE.getValue());
+        runner.setProperty(ConsumeKafka.SCHEMA_CONFLICT_RESOLUTION, 
SchemaConflictResolution.CONTINUE_WITH_MERGED_SCHEMA.getValue());
+        runner.setProperty(ConsumeKafka.AUTO_OFFSET_RESET, 
AutoOffsetReset.EARLIEST.getValue());
+
+        runner.run(1, false, true);
+
+        produceBytes(topic, List.of(
+                serializeAvro(SCHEMA_WITH_ID, Map.of("id", 1L)),
+                serializeAvro(SCHEMA_WITH_NAME, Map.of("name", "Alice"))));
+
+        runUntil(runner, r -> 
!r.getFlowFilesForRelationship(ConsumeKafka.SUCCESS).isEmpty(), 
Duration.ofSeconds(30));
+        runner.run(1, true, false);
+
+        final List<MockFlowFile> successFlowFiles = 
runner.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, successFlowFiles.size());
+
+        final MockFlowFile flowFile = successFlowFiles.getFirst();
+        flowFile.assertAttributeEquals(KafkaFlowFileAttribute.KAFKA_TOPIC, 
topic);
+        flowFile.assertAttributeEquals("record.count", "2");
+
+        final List<GenericRecord> records = new ArrayList<>();
+        final Schema writtenSchema;
+        try (DataFileStream<GenericRecord> stream = new DataFileStream<>(
+                new ByteArrayInputStream(flowFile.toByteArray()), new 
GenericDatumReader<>())) {
+            writtenSchema = stream.getSchema();
+            stream.forEach(records::add);
+        }
+
+        assertEquals(2, records.size());
+
+        final Map<String, Object> firstRecord = new LinkedHashMap<>();
+        firstRecord.put("id", 1L);
+        firstRecord.put("name", null);
+        final Map<String, Object> secondRecord = new LinkedHashMap<>();
+        secondRecord.put("id", null);
+        secondRecord.put("name", "Alice");
+        assertEquals(List.of(firstRecord, secondRecord), 
List.of(toMap(records.get(0)), toMap(records.get(1))));
+
+        assertEquals(List.of("id", "name"), 
writtenSchema.getFields().stream().map(Schema.Field::name).toList());
+        
assertTrue(AvroTypeUtil.isNullable(writtenSchema.getField("id").schema()), "id 
should be nullable in merged schema");
+        
assertTrue(AvroTypeUtil.isNullable(writtenSchema.getField("name").schema()), 
"name should be nullable in merged schema");
+    }
+
+    private static Map<String, Object> toMap(final GenericRecord record) {
+        final Map<String, Object> values = new LinkedHashMap<>();
+        for (final Schema.Field field : record.getSchema().getFields()) {
+            final Object value = record.get(field.name());
+            values.put(field.name(), value instanceof 
org.apache.avro.util.Utf8 utf8 ? utf8.toString() : value);
+        }
+        return values;
+    }
+
+    private static byte[] serializeAvro(final String schemaText, final 
Map<String, Object> values) throws IOException {
+        final Schema schema = new Schema.Parser().parse(schemaText);
+        final GenericRecord record = new GenericData.Record(schema);
+        values.forEach(record::put);
+
+        final ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
+        try (DataFileWriter<GenericRecord> writer = new DataFileWriter<>(new 
GenericDatumWriter<>(schema))) {
+            writer.create(schema, outputStream);
+            writer.append(record);
+        }
+        return outputStream.toByteArray();
+    }
+
+    private void produceBytes(final String topic, final List<byte[]> values)
+            throws ExecutionException, InterruptedException {
+        final Properties properties = new Properties();
+        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 
kafkaContainer.getBootstrapServers());
+        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
ByteArraySerializer.class.getName());
+
+        try (KafkaProducer<String, byte[]> producer = new 
KafkaProducer<>(properties)) {
+            final List<Future<RecordMetadata>> futures = new ArrayList<>();
+            for (final byte[] value : values) {
+                futures.add(producer.send(new ProducerRecord<>(topic, 
FIRST_PARTITION, null, value)));
+            }
+            for (final Future<RecordMetadata> future : futures) {
+                final RecordMetadata metadata = future.get();
+                assertEquals(topic, metadata.topic());
+                assertTrue(metadata.hasOffset());
+            }
+        }
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/ConsumeKafka.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/ConsumeKafka.java
index 77fe7d7fba5..483f5466ff2 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/ConsumeKafka.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/ConsumeKafka.java
@@ -39,9 +39,12 @@ import org.apache.nifi.kafka.processors.common.KafkaUtils;
 import org.apache.nifi.kafka.processors.consumer.OffsetTracker;
 import org.apache.nifi.kafka.processors.consumer.ProcessingStrategy;
 import org.apache.nifi.kafka.processors.consumer.bundle.ByteRecordBundler;
+import 
org.apache.nifi.kafka.processors.consumer.convert.CreateNewFlowFileGrouping;
 import 
org.apache.nifi.kafka.processors.consumer.convert.FlowFileStreamKafkaMessageConverter;
 import 
org.apache.nifi.kafka.processors.consumer.convert.InjectOffsetRecordStreamKafkaMessageConverter;
 import org.apache.nifi.kafka.processors.consumer.convert.KafkaMessageConverter;
+import org.apache.nifi.kafka.processors.consumer.convert.MergeSchemaGrouping;
+import 
org.apache.nifi.kafka.processors.consumer.convert.RecordGroupingStrategy;
 import 
org.apache.nifi.kafka.processors.consumer.convert.RecordStreamKafkaMessageConverter;
 import 
org.apache.nifi.kafka.processors.consumer.convert.WrapperRecordStreamKafkaMessageConverter;
 import org.apache.nifi.kafka.service.api.KafkaConnectionService;
@@ -57,6 +60,7 @@ import org.apache.nifi.kafka.shared.property.HeaderFormat;
 import org.apache.nifi.kafka.shared.property.KeyEncoding;
 import org.apache.nifi.kafka.shared.property.KeyFormat;
 import org.apache.nifi.kafka.shared.property.OutputStrategy;
+import org.apache.nifi.kafka.shared.property.SchemaConflictResolution;
 import org.apache.nifi.logging.ComponentLog;
 import org.apache.nifi.processor.AbstractProcessor;
 import org.apache.nifi.processor.BacklogReportingProcessor;
@@ -97,12 +101,12 @@ import static 
org.apache.nifi.expression.ExpressionLanguageScope.NONE;
 
 @CapabilityDescription("Consumes messages from Apache Kafka Consumer API. "
         + "The complementary NiFi processor for sending messages is 
PublishKafka. The Processor supports consumption of Kafka messages, optionally 
interpreted as NiFi records. "
-        + "Please note that, at this time (in read record mode), the Processor 
assumes that "
-        + "all records that are retrieved from a given partition have the same 
schema. For this mode, if any of the Kafka messages are pulled but cannot be 
parsed or written with the "
+        + "If any of the Kafka messages are pulled but cannot be parsed or 
written with the "
         + "configured Record Reader or Record Writer, the contents of the 
message will be written to a separate FlowFile, and that FlowFile will be 
transferred to the "
         + "'parse.failure' relationship. Otherwise, each FlowFile is sent to 
the 'success' relationship and may contain many individual messages within the 
single FlowFile. "
-        + "A 'record.count' attribute is added to indicate how many messages 
are contained in the FlowFile. No two Kafka messages will be placed into the 
same FlowFile if they "
-        + "have different schemas, or if they have different values for a 
message header that is included by the <Headers to Add as Attributes> property. 
"
+        + "A 'record.count' attribute is added to indicate how many messages 
are contained in the FlowFile. "
+        + "Records are grouped into FlowFiles by topic, partition, and values 
of headers included by the <Headers to Add as Attributes> property. "
+        + "Behavior for conflicting schemas can be configured with the <Schema 
Conflict Resolution> property. "
         + "Kafka Record Header values selected for output are represented 
according to the Header Format property: as text decoded with the configured 
Header Encoding "
         + "character set, or as a lowercase hexadecimal string for binary-safe 
output.")
 @Tags({"Kafka", "Get", "Record", "csv", "avro", "json", "Ingest", "Ingress", 
"Topic", "PubSub", "Consume"})
@@ -274,6 +278,15 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
             .dependsOn(PROCESSING_STRATEGY, ProcessingStrategy.RECORD)
             .build();
 
+    static final PropertyDescriptor SCHEMA_CONFLICT_RESOLUTION = new 
PropertyDescriptor.Builder()
+            .name("Schema Conflict Resolution")
+            .description("Specifies how records with different schemas are 
grouped into FlowFiles.")
+            .required(true)
+            .defaultValue(SchemaConflictResolution.CREATE_NEW_FLOWFILE)
+            .allowableValues(SchemaConflictResolution.class)
+            .dependsOn(PROCESSING_STRATEGY, ProcessingStrategy.RECORD)
+            .build();
+
     static final PropertyDescriptor KEY_ATTRIBUTE_ENCODING = new 
PropertyDescriptor.Builder()
             .name("Key Attribute Encoding")
             .description("Encoding for value of configured FlowFile attribute 
containing Kafka Record Key.")
@@ -348,6 +361,7 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
             RECORD_READER,
             RECORD_WRITER,
             OUTPUT_STRATEGY,
+            SCHEMA_CONFLICT_RESOLUTION,
             KEY_ATTRIBUTE_ENCODING,
             KEY_FORMAT,
             KEY_RECORD_READER,
@@ -364,6 +378,7 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
     private volatile ProcessingStrategy processingStrategy;
     private volatile KeyEncoding keyEncoding;
     private volatile OutputStrategy outputStrategy;
+    private volatile SchemaConflictResolution schemaConflictResolution;
     private volatile KeyFormat keyFormat;
     private volatile boolean commitOffsets;
     private volatile boolean useReader;
@@ -433,6 +448,9 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
                 ? context.getProperty(HEADER_NAME_PREFIX).getValue()
                 : null;
         outputStrategy = processingStrategy == ProcessingStrategy.RECORD ? 
context.getProperty(OUTPUT_STRATEGY).asAllowableValue(OutputStrategy.class) : 
null;
+        schemaConflictResolution = processingStrategy == 
ProcessingStrategy.RECORD
+                ? 
context.getProperty(SCHEMA_CONFLICT_RESOLUTION).asAllowableValue(SchemaConflictResolution.class)
+                : null;
         keyFormat = (outputStrategy == OutputStrategy.USE_WRAPPER || 
outputStrategy == OutputStrategy.INJECT_METADATA)
                 ? 
context.getProperty(KEY_FORMAT).asAllowableValue(KeyFormat.class)
                 : KeyFormat.BYTE_ARRAY;
@@ -867,10 +885,15 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
         final RecordReaderFactory readerFactory = 
context.getProperty(RECORD_READER).asControllerService(RecordReaderFactory.class);
         final RecordSetWriterFactory writerFactory = 
context.getProperty(RECORD_WRITER).asControllerService(RecordSetWriterFactory.class);
 
+        final RecordGroupingStrategy recordGroupingStrategy = switch 
(schemaConflictResolution) {
+            case CONTINUE_WITH_MERGED_SCHEMA -> new 
MergeSchemaGrouping(writerFactory, getLogger(), brokerUri, commitOffsets);
+            case CREATE_NEW_FLOWFILE -> new 
CreateNewFlowFileGrouping(writerFactory, getLogger(), brokerUri, commitOffsets);
+        };
+
         final KafkaMessageConverter converter;
         if (outputStrategy == OutputStrategy.USE_VALUE) {
             converter = new RecordStreamKafkaMessageConverter(readerFactory, 
writerFactory, headerValueConverter, headerNamePattern,
-                    keyEncoding, commitOffsets, offsetTracker, getLogger(), 
brokerUri);
+                    keyEncoding, commitOffsets, offsetTracker, getLogger(), 
brokerUri, recordGroupingStrategy);
         } else if (outputStrategy == OutputStrategy.INJECT_OFFSET) {
             converter = new InjectOffsetRecordStreamKafkaMessageConverter(
                     readerFactory,
@@ -881,14 +904,16 @@ public class ConsumeKafka extends AbstractProcessor 
implements VerifiableProcess
                     commitOffsets,
                     offsetTracker,
                     getLogger(),
-                    brokerUri
+                    brokerUri,
+                    recordGroupingStrategy
             );
         } else {
             final RecordReaderFactory keyReaderFactory = keyFormat == 
KeyFormat.RECORD
                 ? 
context.getProperty(KEY_RECORD_READER).asControllerService(RecordReaderFactory.class)
 : null;
 
             converter = new 
WrapperRecordStreamKafkaMessageConverter(readerFactory, writerFactory, 
keyReaderFactory,
-                headerValueConverter, headerNamePattern, keyFormat, 
keyEncoding, commitOffsets, offsetTracker, getLogger(), brokerUri, 
outputStrategy);
+                headerValueConverter, headerNamePattern, keyFormat, 
keyEncoding, commitOffsets, offsetTracker, getLogger(), brokerUri, 
outputStrategy,
+                recordGroupingStrategy);
         }
 
         converter.toFlowFiles(session, consumerRecords);
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/AbstractRecordStreamKafkaMessageConverter.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/AbstractRecordStreamKafkaMessageConverter.java
index c1f8baff634..bab7e0c93be 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/AbstractRecordStreamKafkaMessageConverter.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/AbstractRecordStreamKafkaMessageConverter.java
@@ -17,41 +17,34 @@
 package org.apache.nifi.kafka.processors.consumer.convert;
 
 import org.apache.nifi.flowfile.FlowFile;
-import org.apache.nifi.flowfile.attributes.CoreAttributes;
 import org.apache.nifi.kafka.processors.ConsumeKafka;
 import org.apache.nifi.kafka.processors.common.HeaderValueConverter;
 import org.apache.nifi.kafka.processors.common.KafkaUtils;
 import org.apache.nifi.kafka.processors.consumer.OffsetTracker;
 import org.apache.nifi.kafka.service.api.record.ByteRecord;
-import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
 import org.apache.nifi.kafka.shared.property.KeyEncoding;
 import org.apache.nifi.logging.ComponentLog;
 import org.apache.nifi.processor.ProcessSession;
-import org.apache.nifi.processor.exception.ProcessException;
 import org.apache.nifi.schema.access.SchemaNotFoundException;
 import org.apache.nifi.serialization.MalformedRecordException;
 import org.apache.nifi.serialization.RecordReader;
 import org.apache.nifi.serialization.RecordReaderFactory;
-import org.apache.nifi.serialization.RecordSetWriter;
 import org.apache.nifi.serialization.RecordSetWriterFactory;
-import org.apache.nifi.serialization.SimpleRecordSchema;
-import org.apache.nifi.serialization.WriteResult;
 import org.apache.nifi.serialization.record.Record;
 import org.apache.nifi.serialization.record.RecordSchema;
 
 import java.io.ByteArrayInputStream;
 import java.io.IOException;
 import java.io.InputStream;
-import java.io.OutputStream;
-import java.util.HashMap;
 import java.util.Iterator;
-import java.util.List;
 import java.util.Map;
-import java.util.concurrent.atomic.AtomicLong;
 import java.util.regex.Pattern;
 
+/**
+ * Shared reader loop, parse-failure handling, and record conversion for 
record-stream converters.
+ * FlowFile grouping and finalization are delegated to a {@link 
RecordGroupingStrategy}.
+ */
 public abstract class AbstractRecordStreamKafkaMessageConverter implements 
KafkaMessageConverter {
-    private static final RecordSchema EMPTY_SCHEMA = new 
SimpleRecordSchema(List.of());
 
     protected final RecordReaderFactory readerFactory;
     protected final RecordSetWriterFactory writerFactory;
@@ -62,8 +55,9 @@ public abstract class 
AbstractRecordStreamKafkaMessageConverter implements Kafka
     protected final OffsetTracker offsetTracker;
     protected final ComponentLog logger;
     protected final String brokerUri;
+    private final RecordGroupingStrategy recordGroupingStrategy;
 
-    public AbstractRecordStreamKafkaMessageConverter(
+    protected AbstractRecordStreamKafkaMessageConverter(
             final RecordReaderFactory readerFactory,
             final RecordSetWriterFactory writerFactory,
             final HeaderValueConverter headerValueConverter,
@@ -72,7 +66,8 @@ public abstract class 
AbstractRecordStreamKafkaMessageConverter implements Kafka
             final boolean commitOffsets,
             final OffsetTracker offsetTracker,
             final ComponentLog logger,
-            final String brokerUri) {
+            final String brokerUri,
+            final RecordGroupingStrategy recordGroupingStrategy) {
         this.readerFactory = readerFactory;
         this.writerFactory = writerFactory;
         this.headerValueConverter = headerValueConverter;
@@ -82,36 +77,32 @@ public abstract class 
AbstractRecordStreamKafkaMessageConverter implements Kafka
         this.offsetTracker = offsetTracker;
         this.logger = logger;
         this.brokerUri = brokerUri;
+        this.recordGroupingStrategy = recordGroupingStrategy;
     }
 
     @Override
     public void toFlowFiles(final ProcessSession session, final 
Iterator<ByteRecord> consumerRecords) {
-        final Map<RecordGroupCriteria, RecordGroup> recordGroups = new 
HashMap<>();
-
         while (consumerRecords.hasNext()) {
             final ByteRecord consumerRecord = consumerRecords.next();
-            final String topic = consumerRecord.getTopic();
-            final int partition = consumerRecord.getPartition();
             final byte[] value = consumerRecord.getValue();
 
-            // shared attribute extraction
             final Map<String, String> attributes = KafkaUtils.toAttributes(
                     consumerRecord, keyEncoding, headerNamePattern, 
headerValueConverter, commitOffsets);
 
-            // hook for subclasses to expose headers (if desired)
-            final Map<String, String> extraAttrs = 
extractHeaders(consumerRecord);
+            final Map<String, String> groupingAttributes = 
extractHeaders(consumerRecord);
 
             try (final InputStream in = new ByteArrayInputStream(value);
                     final RecordReader reader = 
readerFactory.createRecordReader(attributes, in, value.length, logger)) {
 
                 Record record;
                 while ((record = reader.nextRecord()) != null) {
-                    // delegate the actual grouping & writing
-                    processSingleRecord(session, recordGroups, consumerRecord, 
record, attributes, extraAttrs, topic, partition);
+                    final RecordSchema writeSchema = 
getWriteSchema(record.getSchema(), consumerRecord, attributes);
+                    final Record toWrite = convertRecord(consumerRecord, 
record, attributes);
+                    recordGroupingStrategy.addRecord(session, consumerRecord, 
toWrite, writeSchema, attributes, groupingAttributes);
                 }
             } catch (final MalformedRecordException | IOException | 
SchemaNotFoundException e) {
                 logger.debug("Reader or Writer failed to process Kafka Record 
with Topic [{}] Partition [{}] Offset [{}]",
-                             consumerRecord.getTopic(), 
consumerRecord.getPartition(), consumerRecord.getOffset(), e);
+                        consumerRecord.getTopic(), 
consumerRecord.getPartition(), consumerRecord.getOffset(), e);
                 handleParseFailure(session, consumerRecord, attributes, value);
                 offsetTracker.update(consumerRecord);
                 continue;
@@ -122,110 +113,7 @@ public abstract class 
AbstractRecordStreamKafkaMessageConverter implements Kafka
             offsetTracker.update(consumerRecord);
         }
 
-        finishAllGroups(session, recordGroups);
-    }
-
-    private void processSingleRecord(final ProcessSession session,
-            final Map<RecordGroupCriteria, RecordGroup> recordGroups,
-            final ByteRecord consumerRecord,
-            final Record record,
-            final Map<String, String> attributes,
-            final Map<String, String> extraAttrs,
-            final String topic,
-            final int partition) throws Exception {
-        // pick the “bare” schema if the record is null
-        final RecordSchema inputSchema = record == null ? EMPTY_SCHEMA : 
record.getSchema();
-        // let subclass decide how to wrap/transform that into the final schema
-        final RecordSchema writeSchema = getWriteSchema(inputSchema, 
consumerRecord, attributes);
-
-        final RecordGroupCriteria criteria = new 
RecordGroupCriteria(writeSchema, extraAttrs, topic, partition);
-        RecordGroup group = recordGroups.get(criteria);
-        if (group == null) {
-            FlowFile ff = session.create();
-            ff = session.putAllAttributes(ff, Map.of(
-                    KafkaFlowFileAttribute.KAFKA_TOPIC, topic,
-                    KafkaFlowFileAttribute.KAFKA_PARTITION, 
String.valueOf(partition)));
-
-            final OutputStream out = session.write(ff);
-            final RecordSetWriter writer;
-            try {
-                writer = writerFactory.createWriter(logger, writeSchema, out, 
attributes);
-                writer.beginRecordSet();
-            } catch (final Exception ex) {
-                out.close();
-                throw ex;
-            }
-
-            final long offset = consumerRecord.getOffset();
-            final AtomicLong maxOffset = new AtomicLong(offset);
-            final AtomicLong minOffset = new AtomicLong(offset);
-            final AtomicLong minTimestamp = new 
AtomicLong(consumerRecord.getTimestamp());
-            group = new RecordGroup(ff, writer, maxOffset, minOffset, 
minTimestamp);
-            recordGroups.put(criteria, group);
-        } else {
-            final long recordOffset = consumerRecord.getOffset();
-            final AtomicLong maxOffset = group.maxOffset();
-            if (recordOffset > maxOffset.get()) {
-                maxOffset.set(recordOffset);
-            }
-
-            final AtomicLong minOffset = group.minOffset();
-            if (recordOffset < minOffset.get()) {
-                minOffset.set(recordOffset);
-            }
-
-            final long recordTimestamp = consumerRecord.getTimestamp();
-            final AtomicLong minTimestamp = group.minTimestamp();
-            if (recordTimestamp < minTimestamp.get()) {
-                minTimestamp.set(recordTimestamp);
-            }
-        }
-
-        // let subclass convert into the thing to write
-        final Record toWrite = convertRecord(consumerRecord, record, 
attributes);
-        if (toWrite != null) {
-            group.writer().write(toWrite);
-        }
-    }
-
-    private void finishAllGroups(final ProcessSession session, final 
Map<RecordGroupCriteria, RecordGroup> recordGroups) {
-        for (final Map.Entry<RecordGroupCriteria, RecordGroup> e : 
recordGroups.entrySet()) {
-            final RecordGroupCriteria criteria = e.getKey();
-            final RecordGroup group = e.getValue();
-
-            final Map<String, String> resultAttrs = new HashMap<>();
-            final int recordCount;
-            try (final RecordSetWriter writer = group.writer()) {
-                final WriteResult wr = writer.finishRecordSet();
-                resultAttrs.putAll(wr.getAttributes());
-                resultAttrs.put("record.count", 
String.valueOf(wr.getRecordCount()));
-                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_COUNT, 
String.valueOf(wr.getRecordCount()));
-                resultAttrs.put(CoreAttributes.MIME_TYPE.key(), 
writer.getMimeType());
-
-                final long maxOffset = group.maxOffset().get();
-                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_MAX_OFFSET, 
Long.toString(maxOffset));
-
-                final long minOffset = group.minOffset().get();
-                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_OFFSET, 
Long.toString(minOffset));
-
-                final long minTimestamp = group.minTimestamp().get();
-                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_TIMESTAMP, 
Long.toString(minTimestamp));
-
-                // add any extra header‐derived attributes
-                resultAttrs.putAll(criteria.extraAttributes());
-                
resultAttrs.put(KafkaFlowFileAttribute.KAFKA_CONSUMER_OFFSETS_COMMITTED, 
String.valueOf(commitOffsets));
-                recordCount = wr.getRecordCount();
-            } catch (final Exception ex) {
-                throw new ProcessException("Failed to write Kafka records to 
FlowFile", ex);
-            }
-
-            FlowFile ff = group.flowFile();
-            ff = session.putAllAttributes(ff, resultAttrs);
-
-            session.getProvenanceReporter().receive(ff, brokerUri + "/" + 
criteria.topic());
-            session.adjustCounter("Records Received from " + criteria.topic(), 
recordCount, false);
-            session.transfer(ff, ConsumeKafka.SUCCESS);
-        }
+        recordGroupingStrategy.finishAllGroups(session);
     }
 
     protected void handleParseFailure(final ProcessSession session, final 
ByteRecord consumerRecord, final Map<String, String> attributes, final byte[] 
value) {
@@ -246,10 +134,4 @@ public abstract class 
AbstractRecordStreamKafkaMessageConverter implements Kafka
     protected abstract RecordSchema getWriteSchema(RecordSchema inputSchema, 
ByteRecord consumerRecord, Map<String, String> attributes) throws IOException;
 
     protected abstract Record convertRecord(ByteRecord consumerRecord, Record 
record, Map<String, String> attributes) throws IOException;
-
-    private record RecordGroupCriteria(RecordSchema schema, Map<String, 
String> extraAttributes, String topic, int partition) {
-    }
-
-    private record RecordGroup(FlowFile flowFile, RecordSetWriter writer, 
AtomicLong maxOffset, AtomicLong minOffset, AtomicLong minTimestamp) {
-    }
 }
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGrouping.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGrouping.java
new file mode 100644
index 00000000000..adf470a0eb2
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGrouping.java
@@ -0,0 +1,166 @@
+/*
+ * 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.kafka.processors.consumer.convert;
+
+import org.apache.nifi.flowfile.FlowFile;
+import org.apache.nifi.flowfile.attributes.CoreAttributes;
+import org.apache.nifi.kafka.processors.ConsumeKafka;
+import org.apache.nifi.kafka.service.api.record.ByteRecord;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.logging.ComponentLog;
+import org.apache.nifi.processor.ProcessSession;
+import org.apache.nifi.processor.exception.ProcessException;
+import org.apache.nifi.schema.access.SchemaNotFoundException;
+import org.apache.nifi.serialization.RecordSetWriter;
+import org.apache.nifi.serialization.RecordSetWriterFactory;
+import org.apache.nifi.serialization.WriteResult;
+import org.apache.nifi.serialization.record.Record;
+import org.apache.nifi.serialization.record.RecordSchema;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
+
+/**
+ * Create New FlowFile strategy: groups by write schema, topic, partition, and 
grouping attributes,
+ * streaming records into an open writer per group.
+ */
+public class CreateNewFlowFileGrouping implements RecordGroupingStrategy {
+
+    private final RecordSetWriterFactory writerFactory;
+    private final ComponentLog logger;
+    private final String brokerUri;
+    private final boolean commitOffsets;
+    private final Map<RecordGroupCriteria, RecordGroup> recordGroups = new 
HashMap<>();
+
+    public CreateNewFlowFileGrouping(
+            final RecordSetWriterFactory writerFactory,
+            final ComponentLog logger,
+            final String brokerUri,
+            final boolean commitOffsets) {
+        this.writerFactory = writerFactory;
+        this.logger = logger;
+        this.brokerUri = brokerUri;
+        this.commitOffsets = commitOffsets;
+    }
+
+    @Override
+    public void addRecord(
+            final ProcessSession session,
+            final ByteRecord consumerRecord,
+            final Record recordToWrite,
+            final RecordSchema writeSchema,
+            final Map<String, String> attributes,
+            final Map<String, String> groupingAttributes) throws IOException, 
SchemaNotFoundException {
+        final String topic = consumerRecord.getTopic();
+        final int partition = consumerRecord.getPartition();
+
+        final RecordGroupCriteria criteria = new 
RecordGroupCriteria(writeSchema, groupingAttributes, topic, partition);
+        RecordGroup group = recordGroups.get(criteria);
+        if (group == null) {
+            FlowFile ff = session.create();
+            ff = session.putAllAttributes(ff, Map.of(
+                    KafkaFlowFileAttribute.KAFKA_TOPIC, topic,
+                    KafkaFlowFileAttribute.KAFKA_PARTITION, 
String.valueOf(partition)));
+
+            final OutputStream out = session.write(ff);
+            final RecordSetWriter writer;
+            try {
+                writer = writerFactory.createWriter(logger, writeSchema, out, 
attributes);
+                writer.beginRecordSet();
+            } catch (final IOException | SchemaNotFoundException ex) {
+                out.close();
+                throw ex;
+            }
+
+            final long offset = consumerRecord.getOffset();
+            final AtomicLong maxOffset = new AtomicLong(offset);
+            final AtomicLong minOffset = new AtomicLong(offset);
+            final AtomicLong minTimestamp = new 
AtomicLong(consumerRecord.getTimestamp());
+            group = new RecordGroup(ff, writer, maxOffset, minOffset, 
minTimestamp);
+            recordGroups.put(criteria, group);
+        } else {
+            final long recordOffset = consumerRecord.getOffset();
+            final AtomicLong maxOffset = group.maxOffset();
+            if (recordOffset > maxOffset.get()) {
+                maxOffset.set(recordOffset);
+            }
+
+            final AtomicLong minOffset = group.minOffset();
+            if (recordOffset < minOffset.get()) {
+                minOffset.set(recordOffset);
+            }
+
+            final long recordTimestamp = consumerRecord.getTimestamp();
+            final AtomicLong minTimestamp = group.minTimestamp();
+            if (recordTimestamp < minTimestamp.get()) {
+                minTimestamp.set(recordTimestamp);
+            }
+        }
+
+        group.writer().write(recordToWrite);
+    }
+
+    @Override
+    public void finishAllGroups(final ProcessSession session) {
+        for (final Map.Entry<RecordGroupCriteria, RecordGroup> e : 
recordGroups.entrySet()) {
+            final RecordGroupCriteria criteria = e.getKey();
+            final RecordGroup group = e.getValue();
+
+            final Map<String, String> resultAttrs = new HashMap<>();
+            final int recordCount;
+            try (final RecordSetWriter writer = group.writer()) {
+                final WriteResult writeResult = writer.finishRecordSet();
+                resultAttrs.putAll(writeResult.getAttributes());
+                resultAttrs.put("record.count", 
String.valueOf(writeResult.getRecordCount()));
+                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_COUNT, 
String.valueOf(writeResult.getRecordCount()));
+                resultAttrs.put(CoreAttributes.MIME_TYPE.key(), 
writer.getMimeType());
+
+                final long maxOffset = group.maxOffset().get();
+                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_MAX_OFFSET, 
Long.toString(maxOffset));
+
+                final long minOffset = group.minOffset().get();
+                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_OFFSET, 
Long.toString(minOffset));
+
+                final long minTimestamp = group.minTimestamp().get();
+                resultAttrs.put(KafkaFlowFileAttribute.KAFKA_TIMESTAMP, 
Long.toString(minTimestamp));
+
+                resultAttrs.putAll(criteria.groupingAttributes());
+                
resultAttrs.put(KafkaFlowFileAttribute.KAFKA_CONSUMER_OFFSETS_COMMITTED, 
String.valueOf(commitOffsets));
+                recordCount = writeResult.getRecordCount();
+            } catch (final Exception ex) {
+                throw new ProcessException("Failed to write Kafka records to 
FlowFile", ex);
+            }
+
+            FlowFile ff = group.flowFile();
+            ff = session.putAllAttributes(ff, resultAttrs);
+
+            session.getProvenanceReporter().receive(ff, brokerUri + "/" + 
criteria.topic());
+            session.adjustCounter("Records Received from " + criteria.topic(), 
recordCount, false);
+            session.transfer(ff, ConsumeKafka.SUCCESS);
+        }
+        recordGroups.clear();
+    }
+
+    private record RecordGroupCriteria(RecordSchema schema, Map<String, 
String> groupingAttributes, String topic, int partition) {
+    }
+
+    private record RecordGroup(FlowFile flowFile, RecordSetWriter writer, 
AtomicLong maxOffset, AtomicLong minOffset, AtomicLong minTimestamp) {
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/InjectOffsetRecordStreamKafkaMessageConverter.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/InjectOffsetRecordStreamKafkaMessageConverter.java
index a71d9c791c6..5d0f62e3f06 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/InjectOffsetRecordStreamKafkaMessageConverter.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/InjectOffsetRecordStreamKafkaMessageConverter.java
@@ -54,7 +54,8 @@ public class InjectOffsetRecordStreamKafkaMessageConverter 
extends AbstractRecor
             final boolean commitOffsets,
             final OffsetTracker offsetTracker,
             final ComponentLog logger,
-            final String brokerUri
+            final String brokerUri,
+            final RecordGroupingStrategy recordGroupingStrategy
     ) {
         super(
                 readerFactory,
@@ -65,7 +66,8 @@ public class InjectOffsetRecordStreamKafkaMessageConverter 
extends AbstractRecor
                 commitOffsets,
                 offsetTracker,
                 logger,
-                brokerUri
+                brokerUri,
+                recordGroupingStrategy
         );
     }
 
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGrouping.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGrouping.java
new file mode 100644
index 00000000000..7b577211e31
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGrouping.java
@@ -0,0 +1,154 @@
+/*
+ * 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.kafka.processors.consumer.convert;
+
+import org.apache.nifi.flowfile.FlowFile;
+import org.apache.nifi.flowfile.attributes.CoreAttributes;
+import org.apache.nifi.kafka.processors.ConsumeKafka;
+import org.apache.nifi.kafka.service.api.record.ByteRecord;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.logging.ComponentLog;
+import org.apache.nifi.processor.ProcessSession;
+import org.apache.nifi.processor.exception.ProcessException;
+import org.apache.nifi.schema.access.SchemaNotFoundException;
+import org.apache.nifi.serialization.RecordSetWriter;
+import org.apache.nifi.serialization.RecordSetWriterFactory;
+import org.apache.nifi.serialization.WriteResult;
+import org.apache.nifi.serialization.record.Record;
+import org.apache.nifi.serialization.record.RecordSchema;
+import org.apache.nifi.serialization.record.util.DataTypeUtils;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * Continue with Merged Schema strategy: groups by topic, partition, and 
grouping attributes,
+ * merging per-record write schemas and writing once per group.
+ */
+public class MergeSchemaGrouping implements RecordGroupingStrategy {
+
+    private final RecordSetWriterFactory writerFactory;
+    private final ComponentLog logger;
+    private final String brokerUri;
+    private final boolean commitOffsets;
+    private final Map<MergeGroupKey, MergeGroup> mergeGroups = new HashMap<>();
+
+    public MergeSchemaGrouping(
+            final RecordSetWriterFactory writerFactory,
+            final ComponentLog logger,
+            final String brokerUri,
+            final boolean commitOffsets) {
+        this.writerFactory = writerFactory;
+        this.logger = logger;
+        this.brokerUri = brokerUri;
+        this.commitOffsets = commitOffsets;
+    }
+
+    @Override
+    public void addRecord(
+            final ProcessSession session,
+            final ByteRecord consumerRecord,
+            final Record recordToWrite,
+            final RecordSchema writeSchema,
+            final Map<String, String> attributes,
+            final Map<String, String> groupingAttributes) {
+        final MergeGroupKey key = new MergeGroupKey(groupingAttributes, 
consumerRecord.getTopic(), consumerRecord.getPartition());
+        final MergeGroup group = mergeGroups.computeIfAbsent(key, ignored -> 
new MergeGroup(attributes));
+        group.add(recordToWrite, writeSchema, consumerRecord);
+    }
+
+    @Override
+    public void finishAllGroups(final ProcessSession session) {
+        for (final Map.Entry<MergeGroupKey, MergeGroup> entry : 
mergeGroups.entrySet()) {
+            final MergeGroupKey key = entry.getKey();
+            final MergeGroup group = entry.getValue();
+
+            FlowFile flowFile = session.create();
+
+            final RecordSchema schemaToWrite;
+            try {
+                schemaToWrite = writerFactory.getSchema(group.attributes, 
group.mergedWriteSchema);
+            } catch (final SchemaNotFoundException | IOException e) {
+                throw new ProcessException("Failed to determine write schema 
for Kafka records", e);
+            }
+
+            final Map<String, String> flowFileAttributes = new HashMap<>();
+            final AtomicInteger recordCount = new AtomicInteger();
+            flowFile = session.write(flowFile, out -> {
+                try (final RecordSetWriter writer = 
writerFactory.createWriter(logger, schemaToWrite, out, group.attributes)) {
+                    writer.beginRecordSet();
+                    for (final Record record : group.records) {
+                        writer.write(record);
+                    }
+                    final WriteResult writeResult = writer.finishRecordSet();
+                    recordCount.set(writeResult.getRecordCount());
+
+                    flowFileAttributes.putAll(writeResult.getAttributes());
+                    flowFileAttributes.put(CoreAttributes.MIME_TYPE.key(), 
writer.getMimeType());
+                } catch (final SchemaNotFoundException e) {
+                    throw new ProcessException("Failed to write Kafka records 
to FlowFile", e);
+                }
+            });
+
+            flowFileAttributes.put("record.count", 
String.valueOf(recordCount.get()));
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_COUNT, 
String.valueOf(recordCount.get()));
+
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_TOPIC, 
key.topic());
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_PARTITION, 
String.valueOf(key.partition()));
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_MAX_OFFSET, 
Long.toString(group.maxOffset));
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_OFFSET, 
Long.toString(group.minOffset));
+            flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_TIMESTAMP, 
Long.toString(group.minTimestamp));
+            flowFileAttributes.putAll(key.groupingAttributes());
+            
flowFileAttributes.put(KafkaFlowFileAttribute.KAFKA_CONSUMER_OFFSETS_COMMITTED, 
String.valueOf(commitOffsets));
+
+            flowFile = session.putAllAttributes(flowFile, flowFileAttributes);
+
+            session.getProvenanceReporter().receive(flowFile, brokerUri + "/" 
+ key.topic());
+            session.adjustCounter("Records Received from " + key.topic(), 
recordCount.get(), false);
+            session.transfer(flowFile, ConsumeKafka.SUCCESS);
+        }
+        mergeGroups.clear();
+    }
+
+    private record MergeGroupKey(Map<String, String> groupingAttributes, 
String topic, int partition) {
+    }
+
+    private static final class MergeGroup {
+        final Map<String, String> attributes;
+        final List<Record> records = new ArrayList<>();
+        RecordSchema mergedWriteSchema;
+        long maxOffset = Long.MIN_VALUE;
+        long minOffset = Long.MAX_VALUE;
+        long minTimestamp = Long.MAX_VALUE;
+
+        MergeGroup(final Map<String, String> attributes) {
+            this.attributes = attributes;
+        }
+
+        void add(final Record recordToWrite, final RecordSchema writeSchema, 
final ByteRecord consumerRecord) {
+            records.add(recordToWrite);
+            mergedWriteSchema = DataTypeUtils.merge(mergedWriteSchema, 
writeSchema);
+            maxOffset = Math.max(maxOffset, consumerRecord.getOffset());
+            minOffset = Math.min(minOffset, consumerRecord.getOffset());
+            minTimestamp = Math.min(minTimestamp, 
consumerRecord.getTimestamp());
+        }
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordGroupingStrategy.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordGroupingStrategy.java
new file mode 100644
index 00000000000..60d9ce7bf7d
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordGroupingStrategy.java
@@ -0,0 +1,47 @@
+/*
+ * 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.kafka.processors.consumer.convert;
+
+import org.apache.nifi.kafka.service.api.record.ByteRecord;
+import org.apache.nifi.processor.ProcessSession;
+import org.apache.nifi.schema.access.SchemaNotFoundException;
+import org.apache.nifi.serialization.record.Record;
+import org.apache.nifi.serialization.record.RecordSchema;
+
+import java.io.IOException;
+import java.util.Map;
+
+/**
+ * Groups converted Kafka records into FlowFiles according to a Schema 
Conflict Resolution strategy.
+ * <p>
+ * Implementations are stateful: they accumulate open writers or buffered 
records until
+ * {@link #finishAllGroups(ProcessSession)} is called. A new instance must be 
created for each
+ * {@code onTrigger} invocation and must not be reused across calls, so that 
leftover group state
+ * cannot survive an exception or a failed session.
+ */
+public interface RecordGroupingStrategy {
+
+    void addRecord(
+            ProcessSession session,
+            ByteRecord consumerRecord,
+            Record recordToWrite,
+            RecordSchema writeSchema,
+            Map<String, String> attributes,
+            Map<String, String> groupingAttributes) throws IOException, 
SchemaNotFoundException;
+
+    void finishAllGroups(ProcessSession session);
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverter.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverter.java
index 5309db99dc9..27f44b6f821 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverter.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverter.java
@@ -44,8 +44,10 @@ public class RecordStreamKafkaMessageConverter extends 
AbstractRecordStreamKafka
             final boolean commitOffsets,
             final OffsetTracker offsetTracker,
             final ComponentLog logger,
-            final String brokerUri) {
-        super(readerFactory, writerFactory, headerValueConverter, 
headerNamePattern, keyEncoding, commitOffsets, offsetTracker, logger, 
brokerUri);
+            final String brokerUri,
+            final RecordGroupingStrategy recordGroupingStrategy) {
+        super(readerFactory, writerFactory, headerValueConverter, 
headerNamePattern, keyEncoding, commitOffsets, offsetTracker, logger, brokerUri,
+                recordGroupingStrategy);
     }
 
     @Override
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverter.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverter.java
index 8f21d370b34..def21e8108b 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverter.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverter.java
@@ -58,8 +58,10 @@ public class WrapperRecordStreamKafkaMessageConverter 
extends AbstractRecordStre
             final OffsetTracker offsetTracker,
             final ComponentLog logger,
             final String brokerUri,
-            final OutputStrategy outputStrategy) {
-        super(readerFactory, writerFactory, headerValueConverter, 
headerNamePattern, keyEncoding, commitOffsets, offsetTracker, logger, 
brokerUri);
+            final OutputStrategy outputStrategy,
+            final RecordGroupingStrategy recordGroupingStrategy) {
+        super(readerFactory, writerFactory, headerValueConverter, 
headerNamePattern, keyEncoding, commitOffsets, offsetTracker, logger, brokerUri,
+                recordGroupingStrategy);
         this.keyReaderFactory = keyReaderFactory;
         this.keyFormat = keyFormat;
         this.outputStrategy = outputStrategy;
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGroupingTest.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGroupingTest.java
new file mode 100644
index 00000000000..7905c6d1cfd
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/CreateNewFlowFileGroupingTest.java
@@ -0,0 +1,165 @@
+/*
+ * 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.kafka.processors.consumer.convert;
+
+import org.apache.nifi.controller.AbstractControllerService;
+import org.apache.nifi.kafka.processors.ConsumeKafka;
+import org.apache.nifi.kafka.service.api.record.ByteRecord;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.logging.ComponentLog;
+import org.apache.nifi.processor.Processor;
+import org.apache.nifi.reporting.InitializationException;
+import org.apache.nifi.serialization.RecordSetWriter;
+import org.apache.nifi.serialization.RecordSetWriterFactory;
+import org.apache.nifi.serialization.SimpleRecordSchema;
+import org.apache.nifi.serialization.record.MapRecord;
+import org.apache.nifi.serialization.record.MockRecordWriter;
+import org.apache.nifi.serialization.record.Record;
+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.util.MockFlowFile;
+import org.apache.nifi.util.MockProcessSession;
+import org.apache.nifi.util.SharedSessionState;
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CreateNewFlowFileGroupingTest {
+
+    private static final String TOPIC = "topic1";
+    private static final String BROKER_URI = "brokerUri";
+
+    private static final RecordSchema SCHEMA_A = new 
SimpleRecordSchema(List.of(
+            new RecordField("fieldA", RecordFieldType.STRING.getDataType())));
+
+    private static final RecordSchema SCHEMA_B = new 
SimpleRecordSchema(List.of(
+            new RecordField("fieldB", RecordFieldType.INT.getDataType())));
+
+    private static final Record RECORD_A = new MapRecord(SCHEMA_A, 
Map.of("fieldA", "hello"));
+    private static final Record RECORD_B = new MapRecord(SCHEMA_B, 
Map.of("fieldB", 42));
+
+    private final PassThroughSchemaRecordWriter writerFactory = new 
PassThroughSchemaRecordWriter();
+
+    private MockProcessSession session;
+    private ComponentLog logger;
+    private CreateNewFlowFileGrouping grouping;
+
+    @BeforeEach
+    void setUp() throws InitializationException {
+        final TestRunner runner = 
TestRunners.newTestRunner(ConsumeKafka.class);
+        runner.addControllerService("writer", writerFactory);
+        runner.enableControllerService(writerFactory);
+
+        final Processor processor = runner.getProcessor();
+        session = new MockProcessSession(new SharedSessionState(processor, new 
AtomicLong(0)), processor);
+        logger = runner.getLogger();
+        grouping = new CreateNewFlowFileGrouping(writerFactory, logger, 
BROKER_URI, true);
+    }
+
+    @Test
+    void testSameTopicPartitionAndSchemaShareOneFlowFile() throws Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 10, 1000L, 
List.of(), null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 11, 500L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.addRecord(session, second, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, success.size());
+
+        final MockFlowFile flowFile = success.getFirst();
+        assertEquals(TOPIC, 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_TOPIC));
+        assertEquals("0", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION));
+        assertEquals("10", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_OFFSET));
+        assertEquals("11", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_MAX_OFFSET));
+        assertEquals("500", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_TIMESTAMP));
+        assertEquals("2", flowFile.getAttribute("record.count"));
+        assertEquals("true", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_CONSUMER_OFFSETS_COMMITTED));
+    }
+
+    @Test
+    void testDifferentWriteSchemasProduceSeparateFlowFiles() throws Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 1, 1000L, List.of(), 
null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 2, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.addRecord(session, second, RECORD_B, SCHEMA_B, Map.of(), 
Map.of());
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, success.size());
+        assertTrue(success.stream().anyMatch(ff -> 
"1".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_OFFSET))));
+        assertTrue(success.stream().anyMatch(ff -> 
"2".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_OFFSET))));
+    }
+
+    @Test
+    void testDifferentGroupingAttributesProduceSeparateFlowFiles() throws 
Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 1, 1000L, List.of(), 
null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 2, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of("hdr", "a"));
+        grouping.addRecord(session, second, RECORD_A, SCHEMA_A, Map.of(), 
Map.of("hdr", "b"));
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, success.size());
+        assertTrue(success.stream().anyMatch(ff -> 
"a".equals(ff.getAttribute("hdr"))));
+        assertTrue(success.stream().anyMatch(ff -> 
"b".equals(ff.getAttribute("hdr"))));
+    }
+
+    @Test
+    void testDifferentPartitionsProduceSeparateFlowFiles() throws Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 1, 1000L, List.of(), 
null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 1, 2, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.addRecord(session, second, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, success.size());
+        assertTrue(success.stream().anyMatch(ff -> 
"0".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION))));
+        assertTrue(success.stream().anyMatch(ff -> 
"1".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION))));
+    }
+
+    private static final class PassThroughSchemaRecordWriter extends 
AbstractControllerService implements RecordSetWriterFactory {
+        private final MockRecordWriter writer = new MockRecordWriter(null, 
false);
+
+        @Override
+        public RecordSchema getSchema(final Map<String, String> variables, 
final RecordSchema readSchema) {
+            return readSchema;
+        }
+
+        @Override
+        public RecordSetWriter createWriter(final ComponentLog logger, final 
RecordSchema schema, final OutputStream out,
+                final Map<String, String> variables) throws IOException {
+            return writer.createWriter(logger, schema, out, variables);
+        }
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGroupingTest.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGroupingTest.java
new file mode 100644
index 00000000000..c74df340519
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/MergeSchemaGroupingTest.java
@@ -0,0 +1,260 @@
+/*
+ * 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.kafka.processors.consumer.convert;
+
+import org.apache.nifi.controller.AbstractControllerService;
+import org.apache.nifi.kafka.processors.ConsumeKafka;
+import org.apache.nifi.kafka.service.api.record.ByteRecord;
+import org.apache.nifi.kafka.shared.attribute.KafkaFlowFileAttribute;
+import org.apache.nifi.logging.ComponentLog;
+import org.apache.nifi.processor.Processor;
+import org.apache.nifi.reporting.InitializationException;
+import org.apache.nifi.serialization.RecordSetWriter;
+import org.apache.nifi.serialization.RecordSetWriterFactory;
+import org.apache.nifi.serialization.SimpleRecordSchema;
+import org.apache.nifi.serialization.WriteResult;
+import org.apache.nifi.serialization.record.MapRecord;
+import org.apache.nifi.serialization.record.MockRecordWriter;
+import org.apache.nifi.serialization.record.Record;
+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.RecordSet;
+import org.apache.nifi.util.MockFlowFile;
+import org.apache.nifi.util.MockProcessSession;
+import org.apache.nifi.util.SharedSessionState;
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class MergeSchemaGroupingTest {
+
+    private static final String TOPIC = "topic1";
+    private static final String BROKER_URI = "brokerUri";
+
+    private static final RecordSchema SCHEMA_A = new 
SimpleRecordSchema(List.of(
+            new RecordField("fieldA", RecordFieldType.STRING.getDataType())));
+
+    private static final RecordSchema SCHEMA_B = new 
SimpleRecordSchema(List.of(
+            new RecordField("fieldB", RecordFieldType.INT.getDataType())));
+
+    private static final Record RECORD_A = new MapRecord(SCHEMA_A, 
Map.of("fieldA", "hello"));
+    private static final Record RECORD_B = new MapRecord(SCHEMA_B, 
Map.of("fieldB", 42));
+
+    private final PassThroughSchemaRecordWriter writerFactory = new 
PassThroughSchemaRecordWriter();
+
+    private MockProcessSession session;
+    private ComponentLog logger;
+    private MergeSchemaGrouping grouping;
+
+    @BeforeEach
+    void setUp() throws InitializationException {
+        final TestRunner runner = 
TestRunners.newTestRunner(ConsumeKafka.class);
+        runner.addControllerService("writer", writerFactory);
+        runner.enableControllerService(writerFactory);
+
+        final Processor processor = runner.getProcessor();
+        session = new MockProcessSession(new SharedSessionState(processor, new 
AtomicLong(0)), processor);
+        logger = runner.getLogger();
+        grouping = new MergeSchemaGrouping(writerFactory, logger, BROKER_URI, 
true);
+    }
+
+    @Test
+    void testDifferentSchemasInSameTopicPartitionMergeIntoOneFlowFile() throws 
Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 10, 1000L, 
List.of(), null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 11, 500L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.addRecord(session, second, RECORD_B, SCHEMA_B, Map.of(), 
Map.of());
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, success.size());
+
+        final MockFlowFile flowFile = success.getFirst();
+        assertEquals(TOPIC, 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_TOPIC));
+        assertEquals("0", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION));
+        assertEquals("10", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_OFFSET));
+        assertEquals("11", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_MAX_OFFSET));
+        assertEquals("500", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_TIMESTAMP));
+        assertEquals("2", flowFile.getAttribute("record.count"));
+        assertEquals("true", 
flowFile.getAttribute(KafkaFlowFileAttribute.KAFKA_CONSUMER_OFFSETS_COMMITTED));
+
+        assertEquals("hello,\n,42\n", flowFile.getContent());
+    }
+
+    @Test
+    void testDifferentPartitionsProduceSeparateFlowFiles() throws Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 1, 1000L, List.of(), 
null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 1, 2, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of());
+        grouping.addRecord(session, second, RECORD_B, SCHEMA_B, Map.of(), 
Map.of());
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, success.size());
+        assertTrue(success.stream().anyMatch(ff -> 
"0".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION))));
+        assertTrue(success.stream().anyMatch(ff -> 
"1".equals(ff.getAttribute(KafkaFlowFileAttribute.KAFKA_PARTITION))));
+    }
+
+    @Test
+    void testDifferentGroupingAttributesProduceSeparateFlowFiles() throws 
Exception {
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 1, 1000L, List.of(), 
null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 2, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        grouping.addRecord(session, first, RECORD_A, SCHEMA_A, Map.of(), 
Map.of("hdr", "a"));
+        grouping.addRecord(session, second, RECORD_B, SCHEMA_B, Map.of(), 
Map.of("hdr", "b"));
+        grouping.finishAllGroups(session);
+
+        final List<MockFlowFile> success = 
session.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(2, success.size());
+        assertTrue(success.stream().anyMatch(ff -> 
"a".equals(ff.getAttribute("hdr"))));
+        assertTrue(success.stream().anyMatch(ff -> 
"b".equals(ff.getAttribute("hdr"))));
+    }
+
+    @Test
+    void testDisjointNonNullableFieldsBecomeNullableInMergedSchema() throws 
Exception {
+        final RecordSchema nonNullableSchemaA = new SimpleRecordSchema(List.of(
+                new RecordField("fieldA", 
RecordFieldType.STRING.getDataType(), false)));
+        final RecordSchema nonNullableSchemaB = new SimpleRecordSchema(List.of(
+                new RecordField("fieldB", RecordFieldType.INT.getDataType(), 
false)));
+        final Record recordA = new MapRecord(nonNullableSchemaA, 
Map.of("fieldA", "hello"));
+        final Record recordB = new MapRecord(nonNullableSchemaB, 
Map.of("fieldB", 42));
+
+        final AtomicReference<RecordSchema> capturedSchema = new 
AtomicReference<>();
+        final SchemaValidatingRecordWriter validatingWriter = new 
SchemaValidatingRecordWriter(capturedSchema);
+
+        final TestRunner runner = 
TestRunners.newTestRunner(ConsumeKafka.class);
+        runner.addControllerService("validating-writer", validatingWriter);
+        runner.enableControllerService(validatingWriter);
+
+        final Processor processor = runner.getProcessor();
+        final MockProcessSession validatingSession = new MockProcessSession(
+                new SharedSessionState(processor, new AtomicLong(0)), 
processor);
+        final MergeSchemaGrouping validatingGrouping = new MergeSchemaGrouping(
+                validatingWriter, runner.getLogger(), BROKER_URI, true);
+
+        final ByteRecord first = new ByteRecord(TOPIC, 0, 10, 1000L, 
List.of(), null, new byte[0], 0L);
+        final ByteRecord second = new ByteRecord(TOPIC, 0, 11, 2000L, 
List.of(), null, new byte[0], 0L);
+
+        validatingGrouping.addRecord(validatingSession, first, recordA, 
nonNullableSchemaA, Map.of(), Map.of());
+        validatingGrouping.addRecord(validatingSession, second, recordB, 
nonNullableSchemaB, Map.of(), Map.of());
+        validatingGrouping.finishAllGroups(validatingSession);
+
+        final List<MockFlowFile> success = 
validatingSession.getFlowFilesForRelationship(ConsumeKafka.SUCCESS);
+        assertEquals(1, success.size());
+        assertEquals("2", success.getFirst().getAttribute("record.count"));
+        assertEquals("hello,\n,42\n", success.getFirst().getContent());
+
+        final RecordSchema mergedSchema = capturedSchema.get();
+        assertEquals(List.of("fieldA", "fieldB"), 
mergedSchema.getFieldNames());
+        assertTrue(mergedSchema.getField("fieldA").orElseThrow().isNullable());
+        assertTrue(mergedSchema.getField("fieldB").orElseThrow().isNullable());
+    }
+
+    private static final class PassThroughSchemaRecordWriter extends 
AbstractControllerService implements RecordSetWriterFactory {
+        private final MockRecordWriter writer = new MockRecordWriter(null, 
false);
+
+        @Override
+        public RecordSchema getSchema(final Map<String, String> variables, 
final RecordSchema readSchema) {
+            return readSchema;
+        }
+
+        @Override
+        public RecordSetWriter createWriter(final ComponentLog logger, final 
RecordSchema schema, final OutputStream out,
+                final Map<String, String> variables) throws IOException {
+            return writer.createWriter(logger, schema, out, variables);
+        }
+    }
+
+    /**
+     * Captures the write schema and rejects records that omit non-nullable 
fields.
+     */
+    private static final class SchemaValidatingRecordWriter extends 
AbstractControllerService implements RecordSetWriterFactory {
+        private final MockRecordWriter writer = new MockRecordWriter(null, 
false);
+        private final AtomicReference<RecordSchema> capturedSchema;
+
+        private SchemaValidatingRecordWriter(final 
AtomicReference<RecordSchema> capturedSchema) {
+            this.capturedSchema = capturedSchema;
+        }
+
+        @Override
+        public RecordSchema getSchema(final Map<String, String> variables, 
final RecordSchema readSchema) {
+            capturedSchema.set(readSchema);
+            return readSchema;
+        }
+
+        @Override
+        public RecordSetWriter createWriter(final ComponentLog logger, final 
RecordSchema schema, final OutputStream out,
+                final Map<String, String> variables) throws IOException {
+            final RecordSetWriter delegate = writer.createWriter(logger, 
schema, out, variables);
+            return new RecordSetWriter() {
+                @Override
+                public void beginRecordSet() throws IOException {
+                    delegate.beginRecordSet();
+                }
+
+                @Override
+                public WriteResult finishRecordSet() throws IOException {
+                    return delegate.finishRecordSet();
+                }
+
+                @Override
+                public WriteResult write(final Record record) throws 
IOException {
+                    for (final RecordField field : schema.getFields()) {
+                        if (!field.isNullable() && record.getValue(field) == 
null && field.getDefaultValue() == null) {
+                            throw new IOException("Missing required field: " + 
field.getFieldName());
+                        }
+                    }
+                    return delegate.write(record);
+                }
+
+                @Override
+                public WriteResult write(final RecordSet recordSet) throws 
IOException {
+                    return delegate.write(recordSet);
+                }
+
+                @Override
+                public String getMimeType() {
+                    return delegate.getMimeType();
+                }
+
+                @Override
+                public void close() throws IOException {
+                    delegate.close();
+                }
+
+                @Override
+                public void flush() throws IOException {
+                    delegate.flush();
+                }
+            };
+        }
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverterTest.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverterTest.java
index a3aa43be57c..8c544769de9 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverterTest.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/RecordStreamKafkaMessageConverterTest.java
@@ -78,7 +78,8 @@ class RecordStreamKafkaMessageConverterTest {
                 true,
                 offsetTracker,
                 logger,
-                "brokerUri"
+                "brokerUri",
+                new CreateNewFlowFileGrouping(writerFactory, logger, 
"brokerUri", true)
         );
 
         // Create ByteRecords
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverterTest.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverterTest.java
index feaf50ed931..4a7d83ecfa5 100644
--- 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverterTest.java
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/test/java/org/apache/nifi/kafka/processors/consumer/convert/WrapperRecordStreamKafkaMessageConverterTest.java
@@ -83,7 +83,8 @@ class WrapperRecordStreamKafkaMessageConverterTest {
                 offsetTracker,
                 logger,
                 "brokerUri",
-                OutputStrategy.USE_WRAPPER
+                OutputStrategy.USE_WRAPPER,
+                new CreateNewFlowFileGrouping(writerFactory, logger, 
"brokerUri", true)
         );
 
         // Create ByteRecords
diff --git 
a/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-shared/src/main/java/org/apache/nifi/kafka/shared/property/SchemaConflictResolution.java
 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-shared/src/main/java/org/apache/nifi/kafka/shared/property/SchemaConflictResolution.java
new file mode 100644
index 00000000000..2d3f7c8123a
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-shared/src/main/java/org/apache/nifi/kafka/shared/property/SchemaConflictResolution.java
@@ -0,0 +1,55 @@
+/*
+ * 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.kafka.shared.property;
+
+import org.apache.nifi.components.DescribedValue;
+
+/**
+ * Enumeration of supported strategies for resolving schema conflicts when 
batching
+ * Kafka Records into FlowFiles using RECORD processing strategy.
+ */
+public enum SchemaConflictResolution implements DescribedValue {
+    CREATE_NEW_FLOWFILE("Create New FlowFile",
+            "When records have different schemas, a new FlowFile is created 
for each distinct schema."),
+
+    CONTINUE_WITH_MERGED_SCHEMA("Continue with Merged Schema",
+            "When records have different schemas, all schemas within the same 
topic/partition group are merged "
+                    + "into a single write schema so that all records are 
written into a single FlowFile.");
+
+    private final String displayName;
+    private final String description;
+
+    SchemaConflictResolution(final String displayName, final String 
description) {
+        this.displayName = displayName;
+        this.description = description;
+    }
+
+    @Override
+    public String getValue() {
+        return name();
+    }
+
+    @Override
+    public String getDisplayName() {
+        return displayName;
+    }
+
+    @Override
+    public String getDescription() {
+        return description;
+    }
+}

Reply via email to