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]

Reply via email to