github-actions[bot] commented on code in PR #68616:
URL: https://github.com/apache/doris/pull/68616#discussion_r4130609486


##########
fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/deserialize/PostgresDebeziumJsonDeserializer.java:
##########
@@ -101,89 +102,80 @@ 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);
+        }
+
+        Set<String> excludedCols =
+                excludeColumnsCache.getOrDefault(tableId.table(), 
Collections.emptySet());
+        String unsupportedReason =
+                unsupportedChangeReason(stored.getTable(), freshTable, 
excludedCols);
+        if (unsupportedReason != null) {
+            return unsupportedSchemaChange(updatedSchemas, unsupportedReason);
         }
 
         List<Column> added = new ArrayList<>();
         List<String> dropped = new ArrayList<>();
         for (Column col : freshTable.columns()) {
-            if (stored.getTable().columnWithName(col.name()) == null) {
+            if (!excludedCols.contains(col.name())

Review Comment:
   [P1] Check rename ambiguity before excluding columns. With 
exclude_columns=secret, a source table can drop secret, commit that Relation, 
then rename a retained column old to secret. This filter suppresses the ADD of 
secret while the next loop retains DROP old, so the simultaneous ADD/DROP guard 
is bypassed and Doris executes DROP COLUMN old, deleting its historical values 
without the required pause. Compare the full before/after names for rename 
ambiguity first, then apply the exclusion to target DDL.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterJobCommand.java:
##########
@@ -309,7 +301,7 @@ private void checkUnmodifiableProperties(String 
originExecuteSql) throws Analysi
                 break;
             case "cdc_stream":
                 // type, jdbc_url, database, schema, and table identify the 
source and cannot be changed.
-                // snapshot_* are materialized into split metadata on first 
fetch and never re-read.
+                // snapshot_split_key determines persisted split boundaries 
and cannot be changed.

Review Comment:
   [P2] Validate cdc_stream snapshot settings before accepting ALTER. For a 
paused cdc_stream job, setting snapshot_parallelism to 'oops' passes this 
unbound-TVF comparison and is journaled; RESUME then throws in 
JdbcTvfSourceOffsetProvider.ensureInitialized before dispatch, leaving the job 
PENDING and retrying the failure. The same path skips snapshot_split_size 
validation and the retained server_id range check when parallelism increases. 
Run the TVF's integer and merged server_id checks before storing the new SQL.



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