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]