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);
+ }
+ }
}