This is an automated email from the ASF dual-hosted git repository. SvenO3 pushed a commit to branch fix-schema-handling-of-postresql-sink in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 1ae6ce473177a685f4bf1030b223cad8fad7551a Author: Sven Oehler <[email protected]> AuthorDate: Mon Jul 27 16:21:39 2026 +0200 Fix postgresql schema handling --- .../sinks/databases/jvm/jdbcclient/JdbcClient.java | 15 ++++++++++++++- .../sinks/databases/jvm/postgresql/PostgreSql.java | 18 ++++++++++++++++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/jdbcclient/JdbcClient.java b/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/jdbcclient/JdbcClient.java index 1e91b85c9a..9fd62667fe 100644 --- a/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/jdbcclient/JdbcClient.java +++ b/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/jdbcclient/JdbcClient.java @@ -185,8 +185,13 @@ public class JdbcClient { // Database should exist by now so we can establish a connection connection = DriverManager.getConnection(url + databaseName, this.dbDescription.getUsername(), this.dbDescription.getPassword()); + prepareConnection(); this.statementHandler.setStatement(connection.createStatement()); - ResultSet rs = connection.getMetaData().getTables(null, null, this.tableDescription.getName(), null); + ResultSet rs = connection.getMetaData().getTables( + null, + getTableSchemaPattern(), + this.tableDescription.getName(), + null); boolean tableAlreadyExists = rs.next(); rs.close(); if (tableAlreadyExists) { @@ -207,6 +212,14 @@ public class JdbcClient { } } + protected void prepareConnection() throws SQLException { + + } + + protected String getTableSchemaPattern() throws SQLException { + return null; + } + /** * Prepares a statement for the insertion of values or the * diff --git a/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/postgresql/PostgreSql.java b/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/postgresql/PostgreSql.java index 868c94c345..4dacf8579e 100644 --- a/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/postgresql/PostgreSql.java +++ b/streampipes-extensions/streampipes-sinks-databases-jvm/src/main/java/org/apache/streampipes/sinks/databases/jvm/postgresql/PostgreSql.java @@ -26,9 +26,12 @@ import org.apache.streampipes.sinks.databases.jvm.jdbcclient.model.SupportedDbEn import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.sql.SQLException; + public class PostgreSql extends JdbcClient { private static final Logger LOG = LoggerFactory.getLogger(PostgreSql.class); + private static final String DEFAULT_SCHEMA = "public"; private PostgreSqlParameters params; @@ -48,6 +51,21 @@ public class PostgreSql extends JdbcClient { SupportedDbEngines.POSTGRESQL); } + @Override + protected void prepareConnection() throws SQLException { + try (var statement = this.connection.prepareStatement("SELECT current_schema()"); + var resultSet = statement.executeQuery()) { + if (resultSet.next() && resultSet.getString(1) == null) { + this.connection.setSchema(DEFAULT_SCHEMA); + } + } + } + + @Override + protected String getTableSchemaPattern() throws SQLException { + return this.connection.getSchema(); + } + @Override protected void extractTableInformation() {
