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() {
 

Reply via email to