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]