[
https://issues.apache.org/jira/browse/FLINK-40512?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Hongshun Wang updated FLINK-40512:
----------------------------------
Description:
When a PostgreSQL source uses a publication covering all tables
(`publication.autocreate.mode=all_tables` / `CREATE PUBLICATION ... FOR ALL
TABLES`),
the Fluss YAML connector may receive schema events for tables that are not
included
in the configured CDC table list.
was:
When a PostgreSQL source uses a publication covering all tables
(`publication.autocreate.mode=all_tables` / `CREATE PUBLICATION ... FOR ALL
TABLES`),
the Fluss YAML connector may receive schema events for tables that are not
included
in the configured CDC table list.
h3. Actual behavior
PostgreSQL pgoutput sends Relation messages for all relations covered by the
publication. Relation-derived schema records for excluded tables are emitted
downstream and may be converted into schema change events, including
CreateTableEvent, by PostgresPipelineRecordEmitter.
As a result, the Fluss connector subscribes to and processes schemas belonging
to unrelated tables.
h3. Expected behavior
Relation messages should still be consumed internally because they are required
for pgoutput tuple decoding. However, relation-derived schema records for tables
excluded by the configured table filter should not be emitted downstream.
Signal events and schema events required for internal state recovery must remain
unaffected.
h3. Root cause
IncrementalSourceStreamFetcher#shouldEmit applies table and monitored-split
filtering only to data change records.
For every non-data-change record, it currently returns true unconditionally:
{code:java}
// always send the schema change event and signal event
// we need record them to state of Flink
return true;
{code}
PostgreSQL Relation messages are represented as schema records. Therefore,
schemas for unrelated tables bypass taskContext.getTableFilter(), even when
those tables are not part of the configured CDC table list.
PostgresPipelineRecordEmitter#handleSchemaChangeRecord may subsequently infer
and emit schema change events for those tables.
h3. Steps to reproduce
1. Create a PostgreSQL publication covering all tables:
{code:sql}
CREATE PUBLICATION my_publication FOR ALL TABLES;
{code}
2. Configure a PostgreSQL Pipeline YAML source to capture only a subset of the
database tables.
3. Configure Fluss as the sink and enable schema change processing.
4. Produce changes that cause pgoutput to send Relation messages for included
and excluded tables.
5. Observe that schema events, such as CreateTableEvent, for excluded tables
reach the downstream Fluss pipeline.
h3. Relevant code
*
org.apache.flink.cdc.connectors.base.source.reader.external.IncrementalSourceStreamFetcher#shouldEmit
*
org.apache.flink.cdc.connectors.postgres.source.reader.PostgresPipelineRecordEmitter#handleSchemaChangeRecord
> PostgreSQL source emits relation schema events for tables outside the
> configured table list
> -------------------------------------------------------------------------------------------
>
> Key: FLINK-40512
> URL: https://issues.apache.org/jira/browse/FLINK-40512
> Project: Flink
> Issue Type: New Feature
> Affects Versions: cdc-3.6.0
> Reporter: Hongshun Wang
> Priority: Major
> Labels: pull-request-available
> Fix For: cdc-3.7.0
>
>
> When a PostgreSQL source uses a publication covering all tables
> (`publication.autocreate.mode=all_tables` / `CREATE PUBLICATION ... FOR ALL
> TABLES`),
> the Fluss YAML connector may receive schema events for tables that are not
> included
> in the configured CDC table list.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)