goutamadwant opened a new issue, #12546: URL: https://github.com/apache/seatunnel/issues/12546
### Search before asking - [X] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22) and found no similar issues. ### What happened With `schema-changes.enabled = true` and PostgreSQL JDBC driver 42.7.5 or newer, an `ADD COLUMN` is not applied when it arrives with the table's first streamed change after the job starts. The new column's values in that transaction are dropped and the job keeps running. The column is only added when PostgreSQL sends another RELATION message for the table. With pgjdbc 42.4.3 the same steps work. Since pgjdbc 42.7.5, `DatabaseMetaData` returns the database name as `TABLE_CAT`, so Debezium tracks the table as `db.schema.table`. pgoutput RELATION ids have no catalog, so `RelationAwarePostgresSchema#applySchemaChangesForTable` does not find the tracked table and skips the schema-change listener. #10843 fixed the same mismatch for the snapshot lookup. `TABLE_CAT` returned by `getTables`/`getColumns`/`getPrimaryKeys`: | pgjdbc | TABLE_CAT | |---|---| | 42.4.3 - 42.7.4 | `null` | | 42.7.5 - 42.7.13 | `<database>` | **Steps** 1. Table with `REPLICA IDENTITY FULL`, job started with the config below, snapshot finished. 2. `ALTER TABLE public.t1 ADD COLUMN extra text; UPDATE public.t1 SET extra = 'e1' WHERE id = 1; INSERT INTO public.t1 VALUES (100, 'n', 'e100');` (first change on `t1` since the job started) 3. `e1` and `e100` never reach the sink. **Before / After** (PostgreSQL 18.6 and 17.9, Zeta, JDBC sink; values of the new column in the sink) | Scenario | pgjdbc 42.4.3 | pgjdbc 42.7.13, dev | pgjdbc 42.7.13, with fix | |---|---|---|---| | ADD COLUMN, then first change since job start | `e1, e100` | column missing, `e1, e100` lost | `e1, e100` | | Table streamed a change before ADD COLUMN | OK | OK | OK | | Savepoint, ADD COLUMN while stopped, restore | OK | OK | OK | ### SeaTunnel Version dev (46b75fe8b), 3.0.0 ### SeaTunnel Config ```conf env { parallelism = 1, job.mode = "STREAMING", checkpoint.interval = 5000 } source { Postgres-CDC { url = "jdbc:postgresql://localhost:5432/db", username = "postgres", password = "***", database-names = ["db"], schema-names = ["public"], table-names = ["db.public.t1"], slot.name = "st", schema-changes.enabled = true } } sink { Jdbc { url = "jdbc:postgresql://localhost:5432/sink", driver = "org.postgresql.Driver", username = "postgres", password = "***", generate_sink_sql = true, database = "sink", table = "public.${table_name}", schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST", data_save_mode = "APPEND_DATA" } } ``` ### Running Command ```shell bin/seatunnel.sh --config pg_cdc.conf ``` ### Error Exception ```log None. The job stays RUNNING. ``` ### Zeta or Flink or Spark Version Zeta ### Java or Scala Version Java 8, Java 11 ### Screenshots _No response_ ### Are you willing to submit PR? - [X] Yes I am willing to submit a PR! ### Code of Conduct - [X] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
