lvyanquan commented on code in PR #3776:
URL: https://github.com/apache/flink-cdc/pull/3776#discussion_r4236106457


##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/utils/StatementUtils.java:
##########
@@ -58,6 +63,29 @@ public static Object[] queryMinMax(JdbcConnection jdbc, 
TableId tableId, String
                 });
     }
 
+    public static Long queryRowCnt(
+            JdbcConnection jdbc, TableId tableId, String columnName, @Nullable 
String filter)
+            throws SQLException {
+
+        if (filter == null) {
+            return queryApproximateRowCnt(jdbc, tableId);
+        }
+
+        final String cntQuery =
+                String.format("SELECT COUNT(1) FROM %s WHERE (%s)", 
quote(tableId), filter);

Review Comment:
   When a snapshot filter is configured, this changes row-count collection from 
the inexpensive approximate count provided by `SHOW TABLE STATUS` to an exact 
`SELECT COUNT(1) ... WHERE ...`.
   
   Depending on the predicate and available indexes, this query may need to 
scan a large portion of the table. It is executed synchronously during table 
analysis, before snapshot splits can be generated and assigned to parallel 
readers. Together with the preceding filtered `MIN/MAX` query, this may cause 
the source table to be scanned twice before snapshot reading starts.
   
   The row count is only used as a heuristic for calculating the distribution 
factor and dynamic chunk size; it is not required for correctness. Could we 
avoid the exact count when a filter is present—for example, by falling back to 
uneven chunk splitting or continuing to use an approximate estimate? This would 
prevent an expensive count query from blocking snapshot initialization on large 
tables.



##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/config/MySqlSourceConfigFactory.java:
##########
@@ -78,6 +78,7 @@ public class MySqlSourceConfigFactory implements Serializable 
{
     private boolean treatTinyInt1AsBoolean = true;
     private boolean useLegacyJsonFormat = true;
     private boolean assignUnboundedChunkFirst = false;
+    private Map<String, String> snapshotFilters = new HashMap<>();

Review Comment:
   The “first matching filter wins” behavior is not reliably preserved here. 
Although `parseAndValidateSnapshotFilters()` returns a `LinkedHashMap`, 
`MySqlSourceConfigFactory` stores the entries in a `HashMap`, so the original 
YAML/DataStream configuration order is lost before `SnapshotFilterUtils` 
evaluates the patterns.
   
   Additionally, the cache in `SnapshotFilterUtils` uses `Map<String, String>` 
as its key. `Map.equals()` and `hashCode()` are order-insensitive, so two 
configurations containing the same rules in different orders may share the same 
cached selector map even though they have different precedence.
   
   This can cause a table matching multiple patterns to receive a different 
filter from the one configured first, changing the set of snapshot rows.
   
   Could we preserve an ordered representation throughout the configuration 
path and either remove the cache or use an order-sensitive cache key, such as 
an immutable list of entries? It would also be helpful to add tests that:
   
   1. Verify rule order survives the full Factory → SourceConfig path.
   2. Use the same overlapping rules in opposite orders and verify that each 
configuration selects its own first match.



##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/utils/SnapshotFilterUtilsTest.java:
##########
@@ -0,0 +1,107 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.cdc.connectors.mysql.source.utils;
+
+import io.debezium.relational.TableId;
+import org.assertj.core.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+/** Unit test for {@link 
org.apache.flink.cdc.connectors.mysql.source.utils.SnapshotFilterUtils}. */
+public class SnapshotFilterUtilsTest {

Review Comment:
   JUnit 5 does not require public test classes. Please make 
`SnapshotFilterUtilsTest` package-private to follow the project’s testing 
conventions.



##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/MySqlSourceITCase.java:
##########
@@ -496,6 +496,191 @@ void testSnapshotSplitReadingFailCrossCheckpoints(String 
tableName, String chunk
         jobClient.cancel().get();
     }
 
+    @ParameterizedTest
+    @MethodSource("parameters")
+    @SuppressWarnings({"rawtypes", "unchecked"})
+    void testSnapshotFilters(String tableName, String chunkColumnName) throws 
Exception {

Review Comment:
   Could we add a test using SnapshotPhaseHooks to update rows before the high 
watermark and verify the behavior with snapshot backfill both enabled and 
skipped?



##########
flink-cdc-e2e-tests/flink-cdc-pipeline-e2e-tests/src/test/java/org/apache/flink/cdc/pipeline/tests/MysqlE2eITCase.java:
##########
@@ -535,4 +535,54 @@ void testDanglingDropTableEventInBinlog() throws Exception 
{
                 "CreateTableEvent{tableId=%s.products, schema=columns={`id` 
INT NOT NULL,`name` VARCHAR(255) NOT NULL 'flink',`description` 
VARCHAR(512),`weight` FLOAT,`enum_c` STRING 'red',`json_c` STRING,`point_c` 
STRING}, primaryKeys=id, options=()}",
                 "DataChangeEvent{tableId=%s.products, before=[106, hammer, 
16oz carpenter's hammer, 1.0, null, null, null], after=[106, hammer, 18oz 
carpenter hammer, 1.0, null, null, null], op=UPDATE, meta=()}");
     }
+
+    @Test
+    void testSnapshotFilters() throws Exception {
+        String pipelineJob =
+                String.format(
+                        "source:\n"
+                                + "  type: mysql\n"
+                                + "  hostname: %s\n"
+                                + "  port: 3306\n"
+                                + "  username: %s\n"
+                                + "  password: %s\n"
+                                + "  tables: %s.\\.*\n"
+                                + "  server-id: 5400-5404\n"
+                                + "  server-time-zone: UTC\n"
+                                + "  scan.snapshot.filters:\n"
+                                + "    - table: %s.customers\n"
+                                + "      filter: id > 102\n"
+                                + "    - table: %s.products\n"
+                                + "      filter: id < 105\n"
+                                + "\n"
+                                + "sink:\n"
+                                + "  type: values\n"
+                                + "\n"
+                                + "pipeline:\n"
+                                + "  parallelism: %d",
+                        INTER_CONTAINER_MYSQL_ALIAS,
+                        MYSQL_TEST_USER,
+                        MYSQL_TEST_PASSWORD,
+                        mysqlInventoryDatabase.getDatabaseName(),
+                        mysqlInventoryDatabase.getDatabaseName(),
+                        mysqlInventoryDatabase.getDatabaseName(),
+                        parallelism);
+
+        submitPipelineJob(pipelineJob);
+        waitUntilJobRunning(Duration.ofSeconds(30));
+        LOG.info("Pipeline job is running");
+
+        // customers: only id > 102 (id=103, 104)
+        // products: only id < 105 (id=101, 102, 103, 104)
+        validateResult(

Review Comment:
   The new E2E test only calls `validateResult()` with the events expected to 
pass the filters. However, `validateResult()` merely waits for each specified 
event to appear; it does not fail when additional events are emitted.
   
   As a result, this test would still pass if `scan.snapshot.filters` were 
completely ignored and the connector synchronized every row, since all expected 
events would still be present in the full output.
   
   Could we also assert that representative excluded rows are absent—for 
example, `customers.id <= 102` and `products.id >= 105`—or collect and compare 
the complete snapshot result set and event count? That would ensure the test 
actually verifies that filtering is applied.



##########
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/debezium/task/MySqlSnapshotSplitReadTask.java:
##########
@@ -242,12 +247,21 @@ private void createDataEventsForTable(
         long exportStart = clock.currentTimeInMillis();
         LOG.info("Exporting data from split '{}' of table {}", 
snapshotSplit.splitId(), table.id());
 
+        if (snapshotFilter != null) {
+            LOG.info(
+                    "Filter for split '{}' of table {} is: {}",
+                    snapshotSplit.splitId(),
+                    table.id(),
+                    snapshotFilter);
+        }
+
         final String selectSql =
                 StatementUtils.buildSplitScanQuery(
                         snapshotSplit.getTableId(),
                         snapshotSplit.getSplitKeyType(),
                         snapshotSplit.getSplitStart() == null,
-                        snapshotSplit.getSplitEnd() == null);
+                        snapshotSplit.getSplitEnd() == null,
+                        snapshotFilter);

Review Comment:
   The filter is only applied to the snapshot SELECT here, while binlog events 
between the low and high watermarks are later backfilled without evaluating the 
filter. Could we clarify whether the filter is intended to apply only to the 
snapshot query or to the final normalized snapshot? Please document the 
expected behavior and add concurrent DML tests covering rows that enter or 
leave the filter condition during the snapshot phase.



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

Reply via email to