[ 
https://issues.apache.org/jira/browse/FLINK-40321?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18101557#comment-18101557
 ] 

牛一凡 commented on FLINK-40321:
-----------------------------

I would like to work on this issue. Please assign it to me.

> [Flink CDC][Postgres] Table-level DDL metrics are not updated for schema 
> change records
> ---------------------------------------------------------------------------------------
>
>                 Key: FLINK-40321
>                 URL: https://issues.apache.org/jira/browse/FLINK-40321
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>         Environment: Flink:1.20.3
> Flink CDC:3.5.0
> PostgreSQL:12.4
>            Reporter: 牛一凡
>            Priority: Major
>
> *Search before asking*
> I searched in the issues and found no similar issue.
> *Flink CDC version*
> master
> *Database and version*
> PostgreSQL
> *What happened*
> When using the PostgreSQL CDC source with schema changes enabled, a schema 
> change event can trigger a warning in the TaskManager log:
> {code:java}
> WARN org.apache.flink.cdc.connectors.base.source.metrics.SourceReaderMetrics 
> - Failed to update metrics
> org.apache.kafka.connect.errors.DataException: source is not a valid field 
> name
> {code}
> The job itself can continue running, but table-level source metrics are not 
> updated correctly for the DDL record.
> For example, after one schema change:
> {code:java}
> Source__Flink_CDC_Event_Source__postgres.numDDLRecords = 1
> Source__Flink_CDC_Event_Source__postgres.namespace.{namespace}.schema.{schema}.table.{table}.numDDLRecords
>  = 0
> {code}
> The reader-level DDL counter is incremented, but the table-level DDL counter 
> is not.
> *How to reproduce*
> # Start a PostgreSQL CDC pipeline/job with parallelism greater than or equal 
> to 1.
> # Enable schema change synchronization.
> # Synchronize a PostgreSQL table, for example 
> postgres.public.schema_change_test.
> # Execute a DDL on the source table, for example:
> {code:sql}
> ALTER TABLE public.schema_change_test ADD COLUMN add_col_text VARCHAR(100);
> {code}
> # Check the TaskManager log and source metrics.
> *Root cause*
> SourceReaderMetrics.updateRecordCounters() treats the record as a schema 
> change event and increments the reader-level DDL counter first. It then calls 
> SourceRecordUtils.getTableId(record) to update table-level metrics.
> SourceRecordUtils.getTableId() expects record.value().source to contain db, 
> schema, and table.
> However, PostgreSQL uses a custom PostgresSchemaRecord. Its key schema 
> matches io.debezium.connector.postgres-cdc.SchemaChangeKey, so it is 
> recognized as a schema change event, but its value schema is currently a 
> minimal empty struct and does not contain the source field.
> As a result, table-level metrics fail with:
> {code:java}
> source is not a valid field name
> {code}
> *Expected behavior*
> PostgreSQL schema change records should expose the same table identity to 
> source metrics as PostgreSQL data change records.
> For a data record from table postgres.public.schema_change_test, metrics use:
> {code:java}
> namespace.postgres.schema.public.table.schema_change_test
> {code}
> The schema change record should use the same table identity, so the 
> table-level DDL metric should become:
> {code:java}
> Source__Flink_CDC_Event_Source__postgres.namespace.postgres.schema.public.table.schema_change_test.numDDLRecords
>  = 1
> {code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to