diff --git a/database-commons/src/main/java/io/cdap/plugin/db/action/AbstractDBAction.java b/database-commons/src/main/java/io/cdap/plugin/db/action/AbstractDBAction.java index a2be9cbf0..0eaac3148 100644 --- a/database-commons/src/main/java/io/cdap/plugin/db/action/AbstractDBAction.java +++ b/database-commons/src/main/java/io/cdap/plugin/db/action/AbstractDBAction.java @@ -16,12 +16,15 @@ package io.cdap.plugin.db.action; +import io.cdap.cdap.etl.api.FailureCollector; import io.cdap.cdap.etl.api.PipelineConfigurer; import io.cdap.cdap.etl.api.action.Action; import io.cdap.cdap.etl.api.action.ActionContext; +import io.cdap.plugin.common.db.DBErrorDetailsProvider; import io.cdap.plugin.util.DBUtils; import java.sql.Driver; +import java.sql.SQLException; /** * Action that runs a db command. @@ -40,7 +43,18 @@ public AbstractDBAction(QueryConfig config, Boolean enableAutoCommit) { public void run(ActionContext context) throws Exception { Class driverClass = context.loadPluginClass(JDBC_PLUGIN_ID); DBRun executeQuery = new DBRun(config, driverClass, enableAutoCommit); - executeQuery.run(); + try { + executeQuery.run(); + } catch (Exception e) { + if (e instanceof SQLException) { + DBErrorDetailsProvider dbe = new DBErrorDetailsProvider(); + throw dbe.getProgramFailureException((SQLException) e, null); + } + FailureCollector collector = context.getFailureCollector(); + collector.addFailure("Failed to execute query with message: " + e.getMessage(), null) + .withStacktrace(e.getStackTrace()); + collector.getOrThrowException(); + } } @Override