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