diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-paimon/src/main/java/org/apache/flink/cdc/connectors/paimon/sink/v2/bucket/BucketAssignOperator.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-paimon/src/main/java/org/apache/flink/cdc/connectors/paimon/sink/v2/bucket/BucketAssignOperator.java index 07509574c62..358d7c36ce3 100644 --- a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-paimon/src/main/java/org/apache/flink/cdc/connectors/paimon/sink/v2/bucket/BucketAssignOperator.java +++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-paimon/src/main/java/org/apache/flink/cdc/connectors/paimon/sink/v2/bucket/BucketAssignOperator.java @@ -130,7 +130,7 @@ public void processElement(StreamRecord streamRecord) throws Exception { if (event instanceof DataChangeEvent) { DataChangeEvent dataChangeEvent = (DataChangeEvent) event; - if (schemaMaps.containsKey(dataChangeEvent.tableId())) { + if (!schemaMaps.containsKey(dataChangeEvent.tableId())) { Optional schema = schemaEvolutionClient.getLatestEvolvedSchema(dataChangeEvent.tableId()); if (schema.isPresent()) {