diff --git a/debezium-connector-postgres/pom.xml b/debezium-connector-postgres/pom.xml index 733b9b2e6..11f6cede3 100644 --- a/debezium-connector-postgres/pom.xml +++ b/debezium-connector-postgres/pom.xml @@ -40,6 +40,8 @@ 8080 60000 + + v1.0 @@ -51,6 +53,11 @@ org.postgresql postgresql + + org.checkerframework + checker-qual + 3.31.0 + com.google.protobuf protobuf-java @@ -160,6 +167,7 @@ + ${project.artifactId}-polardbo-${polardbo.version}-${project.version} com.github.os72 @@ -321,6 +329,39 @@ + + org.apache.maven.plugins + maven-shade-plugin + 3.2.4 + + + shade-debezium + package + + shade + + + false + + + io.debezium:debezium-api + io.debezium:debezium-core + org.postgresql:postgresql + com.google.protobuf:protobuf-java + + + + + org.postgresql + + io.debezium.connector.postgresql.shaded.org.postgresql + + + + + + + diff --git a/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/PolarDBOConnector.java b/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/PolarDBOConnector.java new file mode 100644 index 000000000..391719892 --- /dev/null +++ b/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/PolarDBOConnector.java @@ -0,0 +1,4 @@ +package io.debezium.connector.postgresql; + +public class PolarDBOConnector extends PostgresConnector { +} diff --git a/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java b/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java index c40a249a0..6b0b6b935 100644 --- a/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java +++ b/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java @@ -20,6 +20,7 @@ import org.apache.kafka.connect.errors.ConnectException; import org.postgresql.core.BaseConnection; +import org.postgresql.core.Oid; // POLAR DIFF import org.postgresql.core.TypeInfo; import org.postgresql.jdbc.PgDatabaseMetaData; import org.slf4j.Logger; @@ -117,6 +118,13 @@ private static Map getLongTypeNames() { private int citextArrayOid = Integer.MIN_VALUE; private int ltreeArrayOid = Integer.MIN_VALUE; + /* POLAR DIFF */ + private boolean mappingDateToTimestamp = false; + public static final int ORADATE = 9002; + public static final int ORADATE_ARRAY_V1 = 9005; + public static final int ORADATE_ARRAY_V2 = 9008; + /* POLAR end */ + public TypeRegistry(PostgresConnection connection) { try { this.connection = connection; @@ -130,6 +138,14 @@ public TypeRegistry(PostgresConnection connection) { } private void addType(PostgresType type) { + /* POLAR DIFF */ + int oid = type.getOid(); + if (oid == ORADATE || oid == ORADATE_ARRAY_V1 || oid == ORADATE_ARRAY_V2) { + mappingDateToTimestamp = true; + return; + } + /* POLAR end */ + oidToType.put(type.getOid(), type); nameToType.put(type.getName(), type); @@ -174,6 +190,15 @@ else if (TYPE_NAME_ISBN.equals(type.getName())) { * @return type associated with the given OID */ public PostgresType get(int oid) { + /* POLAR DIFF */ + if (oid == ORADATE) { + oid = Oid.TIMESTAMP; + } + if (oid == ORADATE_ARRAY_V1 || oid == ORADATE_ARRAY_V2) { + oid = Oid.TIMESTAMP_ARRAY; + } + /* POLAR end */ + PostgresType r = oidToType.get(oid); if (r == null) { r = resolveUnknownType(oid); @@ -202,6 +227,18 @@ public PostgresType get(String name) { name = "int8"; break; } + + /* POLAR DIFF */ + if (mappingDateToTimestamp) { + if (name.equals("date")) { + name = "timestamp"; + } + else if (name.equals("_date")) { + name = "_timestamp"; + } + } + /* POLAR end */ + String[] parts = name.split("\\."); if (parts.length > 1) { name = parts[1]; diff --git a/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java b/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java index 1e5c23cad..cf9b3dabf 100644 --- a/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java +++ b/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java @@ -325,6 +325,16 @@ public QueryExecutor openConnectionImpl(HostSpec[] hostSpecs, Properties info) t } runInitialQueries(queryExecutor, info); + /* POLAR DIFF */ + try { + // Lower versions of walsender will report syntax errors + runPolarInitialQueries(queryExecutor); + } catch (org.postgresql.util.PSQLException e) { + if (!e.getMessage().contains("syntax error")) { + LOGGER.log(Level.WARNING, "runPolarInitialQueries failed, due to ", e); + } + } + /* POLAR END */ // And we're done. return queryExecutor; @@ -932,6 +942,37 @@ private void runInitialQueries(QueryExecutor queryExecutor, Properties info) } } + /* + * POLAR DIFF + * After establishing the connection, set some PolarDB-specific parameters by using + * 'select pg_catalog.set_config ... where name = ...' to avoid errors caused by the + * StartupMessage not recognizing these parameters. + * The old version of the replication process will still raise an error, so it is + * necessary to catch exceptions in the outer layer. + */ + private void runPolarInitialQueries(QueryExecutor queryExecutor) throws SQLException { + /* + * Reset nls_xxx_format like dateStyle to make sys.date\timestamp\timestamptz data + * correct Setting this via SQL instead of the standard StartupMessages (like dateStyle) is to + * ensure compatibility with older versions. + */ + String paramSetSql = + "SELECT pg_catalog.set_config(name,boot_val,'f') FROM pg_catalog.pg_settings WHERE name = '%s' "; + SetupQueryRunner.run(queryExecutor, String.format(paramSetSql, "nls_date_format"), false); + SetupQueryRunner.run( + queryExecutor, String.format(paramSetSql, "nls_timestamp_format"), false); + SetupQueryRunner.run( + queryExecutor, String.format(paramSetSql, "nls_timestamp_tz_format"), false); + /* Set polar_prohibit_rowid_logical on. */ + paramSetSql = + "SELECT pg_catalog.set_config(name,'%s','f') FROM pg_catalog.pg_settings WHERE name = '%s' "; + SetupQueryRunner.run( + queryExecutor, + String.format(paramSetSql, "on", "polar_prohibit_rowid_logical"), + false); + } + /* POLAR END */ + /** * Since PG14 there is GUC_REPORT ParamStatus {@code in_hot_standby} which is set to "on" * when the server is in archive recovery or standby mode. In driver's lingo such server is called diff --git a/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java b/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java index e96fd0844..74b09fd4d 100644 --- a/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java +++ b/debezium-connector-postgres/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java @@ -159,6 +159,12 @@ public class QueryExecutorImpl extends QueryExecutorBase { private final AdaptiveFetchCache adaptiveFetchCache; + /* POLAR DIFF */ + public static final int ORADATE = 9002; + public static final int ORADATE_ARRAY_V1 = 9005; + public static final int ORADATE_ARRAY_V2 = 9008; + /* POLAR end */ + @SuppressWarnings({"assignment", "argument", "method.invocation"}) public QueryExecutorImpl(PGStream pgStream, @@ -2184,6 +2190,7 @@ protected void processResults(ResultHandler handler, int flags, boolean adaptive for (int i = 1; i <= numParams; i++) { int typeOid = pgStream.receiveInteger4(); + typeOid = polarMapTypeOid(typeOid); // POLAR DIFF: map sys.date to timestamp params.setResolvedType(i, typeOid); } @@ -2670,6 +2677,7 @@ private Field[] receiveFields() throws IOException { int typeLength = pgStream.receiveInteger2(); int typeModifier = pgStream.receiveInteger4(); int formatType = pgStream.receiveInteger2(); + typeOid = polarMapTypeOid(typeOid); // POLAR DIFF: map sys.date to timestamp fields[i] = new Field(columnLabel, typeOid, typeLength, typeModifier, tableOid, positionInTable); fields[i].setFormat(formatType); @@ -3033,6 +3041,25 @@ public boolean getIntegerDateTimes() { return integerDateTimes; } + /* POLAR DIFF: map sys.date to timestamp */ + public static int polarMapTypeOid(int typeoid) { + switch (typeoid) { + /* Treat SYS.DATE as TIMESTAMP */ + case ORADATE: + typeoid = Oid.TIMESTAMP; + break; + /* Treat SYS.DATE_ARRAY as TIMESTAMP_ARRAY */ + case ORADATE_ARRAY_V1: + case ORADATE_ARRAY_V2: + typeoid = Oid.TIMESTAMP_ARRAY; + break; + default: + break; + } + return typeoid; + } + /* POLAR end */ + private final Deque pendingParseQueue = new ArrayDeque(); private final Deque pendingBindQueue = new ArrayDeque(); private final Deque pendingExecuteQueue = new ArrayDeque(); diff --git a/debezium-connector-postgres/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java b/debezium-connector-postgres/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java index e5e07384c..b8b89a8dd 100644 --- a/debezium-connector-postgres/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java +++ b/debezium-connector-postgres/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java @@ -13,6 +13,7 @@ import org.postgresql.core.ServerVersion; import org.postgresql.core.Tuple; import org.postgresql.core.TypeInfo; +import org.postgresql.core.v3.QueryExecutorImpl; import org.postgresql.util.ByteConverter; import org.postgresql.util.GT; import org.postgresql.util.JdbcBlackHole; @@ -1230,6 +1231,7 @@ public ResultSet getProcedureColumns(@Nullable String catalog, @Nullable String if ("c".equals(returnTypeType) || ("p".equals(returnTypeType) && argModesArray != null)) { String columnsql = "SELECT a.attname,a.atttypid FROM pg_catalog.pg_attribute a " + " WHERE a.attrelid = " + returnTypeRelid + + " AND a.attname <> 'polar_sys_rowid_attr' AND a.attname <> 'polarSeqRowid' " // POLAR DIFF + " AND NOT a.attisdropped AND a.attnum > 0 ORDER BY a.attnum "; Statement columnstmt = connection.createStatement(); ResultSet columnrs = columnstmt.executeQuery(columnsql); @@ -1568,7 +1570,8 @@ public ResultSet getColumns(@Nullable String catalog, @Nullable String schemaPat + " LEFT JOIN pg_catalog.pg_description dsc ON (c.oid=dsc.objoid AND a.attnum = dsc.objsubid) " + " LEFT JOIN pg_catalog.pg_class dc ON (dc.oid=dsc.classoid AND dc.relname='pg_class') " + " LEFT JOIN pg_catalog.pg_namespace dn ON (dc.relnamespace=dn.oid AND dn.nspname='pg_catalog') " - + " WHERE c.relkind in ('r','p','v','f','m') and a.attnum > 0 AND NOT a.attisdropped "; + + " WHERE c.relkind in ('r','p','v','f','m') and a.attnum > 0 AND NOT a.attisdropped " + + " AND a.attname <> 'polar_sys_rowid_attr' AND a.attname <> 'polarSeqRowid' "; // POLAR DIFF if (schemaPattern != null && !schemaPattern.isEmpty()) { sql += " AND n.nspname LIKE " + escapeQuotes(schemaPattern); @@ -1590,6 +1593,7 @@ public ResultSet getColumns(@Nullable String catalog, @Nullable String schemaPat byte[] @Nullable [] tuple = new byte[numberOfFields][]; int typeOid = (int) rs.getLong("atttypid"); int typeMod = rs.getInt("atttypmod"); + typeOid = QueryExecutorImpl.polarMapTypeOid(typeOid); // POLAR DIFF tuple[0] = null; // Catalog name, not supported tuple[1] = rs.getBytes("nspname"); // Schema @@ -1740,7 +1744,8 @@ public ResultSet getColumnPrivileges(@Nullable String catalog, @Nullable String + " AND c.relowner = r.oid " + " AND c.oid = a.attrelid " + " AND c.relkind = 'r' " - + " AND a.attnum > 0 AND NOT a.attisdropped "; + + " AND a.attnum > 0 AND NOT a.attisdropped " + + " AND a.attname <> 'polar_sys_rowid_attr' AND a.attname <> 'polarSeqRowid' "; // POLAR DIFF if (schema != null && !schema.isEmpty()) { sql += " AND n.nspname = " + escapeQuotes(schema); @@ -2169,6 +2174,12 @@ public ResultSet getPrimaryKeys(@Nullable String catalog, @Nullable String schem sql += " AND ct.relname = " + escapeQuotes(table); } + /* POLAR DIFF: support rowid and global index */ + sql += " AND ci.relname <> 'polar_rowid_' || ct.oid || '_index' "; + sql += " AND ci.relname <> 'pg_oid_' || ct.oid || '_index' "; + sql += " AND ci.relname <> 'Polar_Rowid_' || ct.oid || '_index' "; + sql += " AND NOT ('global_index=true' = ANY(ci.reloptions) AND a.attname = 'tableoid') "; + /* POLAR END */ sql += " AND i.indisprimary "; sql = "SELECT " + " result.TABLE_CAT, " @@ -2521,6 +2532,10 @@ public ResultSet getIndexInfo( + " pg_catalog.pg_get_expr(i.indpred, i.indrelid) AS FILTER_CONDITION, " + " ci.oid AS CI_OID, " + " i.indoption AS I_INDOPTION, " + /* POLAR DIFF: support global index */ + + " ('global_index=true' = ANY(ci.reloptions)) AS IS_GLOBAL_INDEX, " + + " (information_schema._pg_expandarray(i.indkey)).x AS COLUMN_ATTNUM, " + /* POLAR END */ + (connection.haveMinimumServerVersion(ServerVersion.v9_6) ? " am.amname AS AM_NAME " : " am.amcanorder AS AM_CANORDER ") + "FROM pg_catalog.pg_class ct " + " JOIN pg_catalog.pg_namespace n ON (ct.relnamespace = n.oid) " @@ -2539,6 +2554,11 @@ public ResultSet getIndexInfo( sql += " AND i.indisunique "; } + /* POLAR DIFF: support rowid */ + sql += " AND ci.relname <> 'polar_rowid_' || ct.oid || '_index' "; + sql += " AND ci.relname <> 'pg_oid_' || ct.oid || '_index' "; + sql += " AND ci.relname <> 'Polar_Rowid_' || ct.oid || '_index' "; + /* POLAR END */ sql = "SELECT " + " tmp.TABLE_CAT, " + " tmp.TABLE_SCHEM, " @@ -2570,6 +2590,7 @@ public ResultSet getIndexInfo( + "FROM (" + sql + ") AS tmp"; + sql += " WHERE NOT (tmp.IS_GLOBAL_INDEX AND tmp.COLUMN_ATTNUM = -6) "; /* POLAR DIFF: support global index */ } else { String select; String from; @@ -3041,6 +3062,7 @@ public ResultSet getFunctionColumns(@Nullable String catalog, @Nullable String s if ("c".equals(returnTypeType) || ("p".equals(returnTypeType) && argModesArray != null)) { String columnsql = "SELECT a.attname,a.atttypid FROM pg_catalog.pg_attribute a " + " WHERE a.attrelid = " + returnTypeRelid + + " AND a.attname <> 'polar_sys_rowid_attr' AND a.attname <> 'polarSeqRowid' " // POLAR DIFF + " AND NOT a.attisdropped AND a.attnum > 0 ORDER BY a.attnum "; Statement columnstmt = connection.createStatement(); ResultSet columnrs = columnstmt.executeQuery(columnsql); diff --git a/debezium-connector-postgres/src/main/resources/META-INF/services/org.apache.kafka.connect.source.SourceConnector b/debezium-connector-postgres/src/main/resources/META-INF/services/org.apache.kafka.connect.source.SourceConnector index 58de09f1b..c72d974dd 100644 --- a/debezium-connector-postgres/src/main/resources/META-INF/services/org.apache.kafka.connect.source.SourceConnector +++ b/debezium-connector-postgres/src/main/resources/META-INF/services/org.apache.kafka.connect.source.SourceConnector @@ -1 +1,2 @@ -io.debezium.connector.postgresql.PostgresConnector \ No newline at end of file +io.debezium.connector.postgresql.PostgresConnector +io.debezium.connector.postgresql.PolarDBOConnector \ No newline at end of file