Copilot commented on code in PR #68616:
URL: https://github.com/apache/doris/pull/68616#discussion_r4131562831


##########
fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/deserialize/PostgresDebeziumJsonDeserializer.java:
##########
@@ -101,26 +102,51 @@ private DeserializeResult handleSchemaChangeEvent(
         TableId tableId = freshTable.id();
         TableChanges.TableChange stored = tableSchemas != null ? 
tableSchemas.get(tableId) : null;
         LOG.info(
-                "[SCHEMA-CHANGE] Postgres deserializer received schema change, 
table={}, baselineSchemas={}, hasStoredSchema={}",
+                "[SCHEMA-CHANGE] Postgres deserializer received schema change, 
table={},"
+                        + " baselineSchemas={}, hasStoredSchema={}",
                 tableId.identifier(),
                 tableSchemas == null ? 0 : tableSchemas.size(),
                 stored != null && stored.getTable() != null);
 
-        // changeType is not consumed inside cdc_client — downstream only 
reads getTable() and
-        // serializeTableSchemas does not persist it — so ALTER is used 
uniformly, including for the
-        // first-time baseline below (which is semantically a CREATE).
+        // Relation events use ALTER uniformly, including when establishing 
the initial baseline.
         TableChanges.TableChange freshChange =
                 new 
TableChanges.TableChange(TableChanges.TableChangeType.ALTER, freshTable);
         Map<TableId, TableChanges.TableChange> updatedSchemas = new 
HashMap<>();
+        updatedSchemas.put(tableId, freshChange);
 
         // No baseline yet: adopt the fresh schema as baseline, no DDL.
         if (stored == null || stored.getTable() == null) {
             LOG.info(
-                    "[SCHEMA-CHANGE] Table {}: no baseline, adopting fresh 
schema as baseline (no DDL)",
+                    "[SCHEMA-CHANGE] Table {}: no baseline, adopting fresh 
schema as baseline (no"
+                            + " DDL)",
+                    tableId.identifier());
+            return DeserializeResult.schemaChange(Collections.emptyList(), 
updatedSchemas);
+        }
+
+        // Debezium equality does not compare native type IDs.
+        if (stored.getTable().equals(freshTable)
+                && stored.getTable().columns().stream()
+                        .allMatch(
+                                column ->
+                                        column.nativeType()
+                                                == freshTable
+                                                        
.columnWithName(column.name())
+                                                        .nativeType())) {
+            return DeserializeResult.empty();
+        }
+        if (isSchemaChangeIgnored(context)) {
+            LOG.info(
+                    "[SCHEMA-CHANGE-IGNORED] Postgres target DDL skipped for 
table {}",
                     tableId.identifier());
-            updatedSchemas.put(tableId, freshChange);
-            return DeserializeResult.schemaChange(
-                    Collections.emptyList(), updatedSchemas, 
Collections.emptyList());
+            return DeserializeResult.schemaChange(Collections.emptyList(), 
updatedSchemas);

Review Comment:
   `ignore` returns before `unsupportedChangeReason`, which is where 
primary-key changes are rejected. As a result, changing a PostgreSQL primary 
key in ignore mode silently advances the source baseline even though the target 
key is unchanged, risking incorrect update/delete behavior. Check primary-key 
equality before this early return.



##########
fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/deserialize/MySqlDebeziumJsonDeserializer.java:
##########
@@ -112,8 +113,29 @@ private DeserializeResult handleSchemaChangeEvent(
                                 + "baseline without emitting Doris DDL. DDL: 
{}",
                         tableId.identifier(),
                         ddl);
-                return DeserializeResult.schemaChange(
-                        Collections.emptyList(), freshSchemas, 
Collections.emptyList());
+                return DeserializeResult.schemaChange(Collections.emptyList(), 
freshSchemas);
+            }
+        }
+
+        if (isSchemaChangeIgnored(context)) {
+            LOG.info(
+                    "[SCHEMA-CHANGE-IGNORED] MySQL target DDL skipped for 
tables {}",
+                    freshSchemas.keySet());
+            return DeserializeResult.schemaChange(Collections.emptyList(), 
freshSchemas);
+        }
+
+        for (Map.Entry<TableId, TableChanges.TableChange> entry : 
freshSchemas.entrySet()) {
+            if (entry.getValue().getType() != 
TableChanges.TableChangeType.ALTER) {
+                continue;
+            }
+            TableId tableId = entry.getKey();
+            if (!tableSchemas
+                    .get(tableId)
+                    .getTable()
+                    .primaryKeyColumnNames()
+                    
.equals(entry.getValue().getTable().primaryKeyColumnNames())) {
+                return unsupportedSchemaChange(
+                        record, ddl, "Primary key changes are not supported", 
freshSchemas);
             }
         }

Review Comment:
   `ignore` returns before the primary-key comparison below, so a source 
primary-key change is silently accepted and the baseline advances. That can 
make subsequent updates/deletes use keys that no longer match the Doris table. 
Perform the primary-key check first; only skip column DDL and other 
column-level checks in ignore mode.



-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to