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]