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

exceptionfactory pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git


The following commit(s) were added to refs/heads/main by this push:
     new 732f77b8041 NIFI-16227 GenerateTableFetch can leave incoming FlowFiles 
unacknowledged after processing failures (#11563)
732f77b8041 is described below

commit 732f77b8041b220559016805c74741cfc303c144
Author: Pierre Villard <[email protected]>
AuthorDate: Fri Aug 21 16:48:46 2026 +0200

    NIFI-16227 GenerateTableFetch can leave incoming FlowFiles unacknowledged 
after processing failures (#11563)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../processors/standard/GenerateTableFetch.java    | 53 +++++++++--------
 .../standard/TestGenerateTableFetch.java           | 66 ++++++++++++++++++++++
 2 files changed, 94 insertions(+), 25 deletions(-)

diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GenerateTableFetch.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GenerateTableFetch.java
index 93f6c488da4..58b0dd0a087 100644
--- 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GenerateTableFetch.java
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/GenerateTableFetch.java
@@ -270,37 +270,37 @@ public class GenerateTableFetch extends 
AbstractDatabaseFetchProcessor {
                 return;
             }
         }
-        maxValueProperties = getDefaultMaxValueProperties(context, 
fileToProcess);
-
         final ComponentLog logger = getLogger();
+        try {
+            maxValueProperties = getDefaultMaxValueProperties(context, 
fileToProcess);
 
-        final DBCPService dbcpService = 
context.getProperty(DBCP_SERVICE).asControllerService(DBCPService.class);
-        final DatabaseDialectService databaseDialectService = 
getDatabaseDialectService(context);
-        final String databaseType = context.getProperty(DB_TYPE).getValue();
+            final DBCPService dbcpService = 
context.getProperty(DBCP_SERVICE).asControllerService(DBCPService.class);
+            final DatabaseDialectService databaseDialectService = 
getDatabaseDialectService(context);
+            final String databaseType = 
context.getProperty(DB_TYPE).getValue();
 
-        final String tableName = 
context.getProperty(TABLE_NAME).evaluateAttributeExpressions(fileToProcess).getValue();
-        final String columnNames = 
context.getProperty(COLUMN_NAMES).evaluateAttributeExpressions(fileToProcess).getValue();
-        final String maxValueColumnNames = 
context.getProperty(MAX_VALUE_COLUMN_NAMES).evaluateAttributeExpressions(fileToProcess).getValue();
-        final int partitionSize = 
context.getProperty(PARTITION_SIZE).evaluateAttributeExpressions(fileToProcess).asInteger();
-        final String columnForPartitioning = 
context.getProperty(COLUMN_FOR_VALUE_PARTITIONING).evaluateAttributeExpressions(fileToProcess).getValue();
-        final boolean useColumnValsForPaging = 
!StringUtils.isEmpty(columnForPartitioning);
-        final String customWhereClause = 
context.getProperty(WHERE_CLAUSE).evaluateAttributeExpressions(fileToProcess).getValue();
-        final String customOrderByColumn = 
context.getProperty(CUSTOM_ORDERBY_COLUMN).evaluateAttributeExpressions(fileToProcess).getValue();
-        final boolean outputEmptyFlowFileOnZeroResults = 
context.getProperty(OUTPUT_EMPTY_FLOWFILE_ON_ZERO_RESULTS).asBoolean();
+            final String tableName = 
context.getProperty(TABLE_NAME).evaluateAttributeExpressions(fileToProcess).getValue();
+            final String columnNames = 
context.getProperty(COLUMN_NAMES).evaluateAttributeExpressions(fileToProcess).getValue();
+            final String maxValueColumnNames = 
context.getProperty(MAX_VALUE_COLUMN_NAMES).evaluateAttributeExpressions(fileToProcess).getValue();
+            final int partitionSize = 
context.getProperty(PARTITION_SIZE).evaluateAttributeExpressions(fileToProcess).asInteger();
+            final String columnForPartitioning = 
context.getProperty(COLUMN_FOR_VALUE_PARTITIONING).evaluateAttributeExpressions(fileToProcess).getValue();
+            final boolean useColumnValsForPaging = 
!StringUtils.isEmpty(columnForPartitioning);
+            final String customWhereClause = 
context.getProperty(WHERE_CLAUSE).evaluateAttributeExpressions(fileToProcess).getValue();
+            final String customOrderByColumn = 
context.getProperty(CUSTOM_ORDERBY_COLUMN).evaluateAttributeExpressions(fileToProcess).getValue();
+            final boolean outputEmptyFlowFileOnZeroResults = 
context.getProperty(OUTPUT_EMPTY_FLOWFILE_ON_ZERO_RESULTS).asBoolean();
 
-        final StateMap stateMap;
-        FlowFile finalFileToProcess = fileToProcess;
+            final StateMap stateMap;
+            FlowFile finalFileToProcess = fileToProcess;
 
-        try {
-            stateMap = session.getState(Scope.CLUSTER);
-        } catch (final IOException ioe) {
-            logger.error("Failed to retrieve observed maximum values from the 
State Manager. Will not perform "
-                    + "query until this is accomplished.", ioe);
-            context.yield();
-            return;
-        }
+            try {
+                stateMap = session.getState(Scope.CLUSTER);
+            } catch (final IOException ioe) {
+                logger.error("Failed to retrieve observed maximum values from 
the State Manager. Will not perform "
+                        + "query until this is accomplished.", ioe);
+                session.rollback();
+                context.yield();
+                return;
+            }
 
-        try {
             // Make a mutable copy of the current state property map. This 
will be updated by the result row callback, and eventually
             // set as the current state map (after the session has been 
committed)
             final Map<String, String> statePropertyMap = new 
HashMap<>(stateMap.toMap());
@@ -584,6 +584,9 @@ public class GenerateTableFetch extends 
AbstractDatabaseFetchProcessor {
             logger.error("Error during processing: {}", t.getMessage(), t);
             session.rollback();
             context.yield();
+        } catch (final RuntimeException e) {
+            session.rollback();
+            throw e;
         }
     }
 
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestGenerateTableFetch.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestGenerateTableFetch.java
index 1f33a5f0d95..030f24a6ae6 100644
--- 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestGenerateTableFetch.java
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestGenerateTableFetch.java
@@ -18,6 +18,12 @@ package org.apache.nifi.processors.standard;
 
 import org.apache.nifi.components.state.Scope;
 import org.apache.nifi.components.state.StateManager;
+import org.apache.nifi.controller.AbstractControllerService;
+import org.apache.nifi.database.dialect.service.api.DatabaseDialectService;
+import org.apache.nifi.database.dialect.service.api.StatementRequest;
+import org.apache.nifi.database.dialect.service.api.StatementResponse;
+import org.apache.nifi.database.dialect.service.api.StatementType;
+import org.apache.nifi.reporting.InitializationException;
 import org.apache.nifi.util.MockFlowFile;
 import org.apache.nifi.util.MockProcessSession;
 import org.apache.nifi.util.MockSessionFactory;
@@ -46,6 +52,7 @@ import static 
org.apache.nifi.processors.standard.AbstractDatabaseFetchProcessor
 import static 
org.apache.nifi.processors.standard.AbstractDatabaseFetchProcessor.WHERE_CLAUSE;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 
 class TestGenerateTableFetch extends AbstractDatabaseConnectionServiceTest {
 
@@ -1118,6 +1125,53 @@ class TestGenerateTableFetch extends 
AbstractDatabaseConnectionServiceTest {
         assertEquals(expectedRenamed, 
propertyMigrationResult.getPropertiesRenamed());
     }
 
+    @Test
+    void testUncheckedDialectFailureRollsBackIncomingFlowFile() throws 
InitializationException {
+        final DatabaseDialectService dialectService = new 
FailingDatabaseDialectService();
+        runner.addControllerService("failing-dialect", dialectService);
+        runner.enableControllerService(dialectService);
+        runner.setProperty(DB_TYPE, "Database Dialect Service");
+        runner.setProperty(GenerateTableFetch.DATABASE_DIALECT_SERVICE, 
"failing-dialect");
+        runner.setProperty(GenerateTableFetch.TABLE_NAME, "${tableName}");
+        runner.setIncomingConnection(true);
+        runner.enqueue("", Map.of("tableName", "TEST_QUERY_DB_TABLE"));
+
+        final AssertionError error = assertThrows(AssertionError.class, 
runner::run);
+        assertEquals(IllegalArgumentException.class, 
error.getCause().getClass());
+
+        final MockProcessSession session = ((MockSessionFactory) 
runner.getProcessSessionFactory()).getCreatedSessions().iterator().next();
+        session.assertRolledBack();
+        runner.assertQueueNotEmpty();
+    }
+
+    @Test
+    void testInvalidFlowFilePropertyRollsBackIncomingFlowFile() {
+        runner.setProperty(GenerateTableFetch.TABLE_NAME, "${tableName}");
+        runner.setProperty(GenerateTableFetch.PARTITION_SIZE, "${partSize}");
+        runner.setIncomingConnection(true);
+        runner.enqueue("", Map.of("tableName", "TEST_QUERY_DB_TABLE", 
"partSize", "invalid"));
+
+        assertThrows(AssertionError.class, runner::run);
+
+        final MockProcessSession session = ((MockSessionFactory) 
runner.getProcessSessionFactory()).getCreatedSessions().iterator().next();
+        session.assertRolledBack();
+        runner.assertQueueNotEmpty();
+    }
+
+    @Test
+    void testStateReadFailureRollsBackIncomingFlowFile() {
+        runner.setProperty(GenerateTableFetch.TABLE_NAME, "${tableName}");
+        runner.setIncomingConnection(true);
+        runner.enqueue("", Map.of("tableName", "TEST_QUERY_DB_TABLE"));
+        runner.getStateManager().setFailOnStateGet(Scope.CLUSTER, true);
+
+        runner.run();
+
+        final MockProcessSession session = ((MockSessionFactory) 
runner.getProcessSessionFactory()).getCreatedSessions().iterator().next();
+        session.assertRolledBack();
+        runner.assertQueueNotEmpty();
+    }
+
     private void assertResultsFound(final String query, final int results) 
throws SQLException {
         int resultsFound = 0;
         try (
@@ -1131,4 +1185,16 @@ class TestGenerateTableFetch extends 
AbstractDatabaseConnectionServiceTest {
         }
         assertEquals(results, resultsFound);
     }
+
+    private static class FailingDatabaseDialectService extends 
AbstractControllerService implements DatabaseDialectService {
+        @Override
+        public StatementResponse getStatement(final StatementRequest 
statementRequest) {
+            throw new IllegalArgumentException("Order By is required when 
paging is specified");
+        }
+
+        @Override
+        public Set<StatementType> getSupportedStatementTypes() {
+            return Set.of(StatementType.SELECT);
+        }
+    }
 }

Reply via email to