This is an automated email from the ASF dual-hosted git repository.

jiabaosun pushed a commit to branch v2.0
in repository https://gitbox.apache.org/repos/asf/flink-connector-mongodb.git


The following commit(s) were added to refs/heads/v2.0 by this push:
     new e9d484b  [FLINK-37255] Fix unable to configure 
scan.partition.record-size option
e9d484b is described below

commit e9d484b1e040de9d9c4f667ac941a2081348d435
Author: yuxiqian <[email protected]>
AuthorDate: Wed Feb 5 19:10:12 2025 +0800

    [FLINK-37255] Fix unable to configure scan.partition.record-size option
---
 .../mongodb/source/config/MongoReadOptions.java    | 13 ++++++--
 .../mongodb/table/MongoDynamicTableFactory.java    |  2 ++
 .../mongodb/source/MongoSourceITCase.java          |  3 +-
 .../table/MongoDynamicTableFactoryTest.java        | 36 ++++++++++++++++++++++
 4 files changed, 51 insertions(+), 3 deletions(-)

diff --git 
a/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/source/config/MongoReadOptions.java
 
b/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/source/config/MongoReadOptions.java
index 41590e1..de41b64 100644
--- 
a/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/source/config/MongoReadOptions.java
+++ 
b/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/source/config/MongoReadOptions.java
@@ -104,12 +104,18 @@ public class MongoReadOptions implements Serializable {
         return noCursorTimeout == that.noCursorTimeout
                 && partitionStrategy == that.partitionStrategy
                 && samplesPerPartition == that.samplesPerPartition
-                && Objects.equals(partitionSize, that.partitionSize);
+                && Objects.equals(partitionSize, that.partitionSize)
+                && Objects.equals(partitionRecordSize, 
that.partitionRecordSize);
     }
 
     @Override
     public int hashCode() {
-        return Objects.hash(noCursorTimeout, partitionStrategy, partitionSize, 
samplesPerPartition);
+        return Objects.hash(
+                noCursorTimeout,
+                partitionStrategy,
+                partitionSize,
+                samplesPerPartition,
+                partitionRecordSize);
     }
 
     public static MongoReadOptionsBuilder builder() {
@@ -219,6 +225,9 @@ public class MongoReadOptions implements Serializable {
          */
         public MongoReadOptionsBuilder setPartitionRecordSize(
                 @Nullable Integer partitionRecordSize) {
+            checkArgument(
+                    partitionRecordSize == null || partitionRecordSize > 0,
+                    "The record size per partition must be larger than 0.");
             this.partitionRecordSize = partitionRecordSize;
             return this;
         }
diff --git 
a/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactory.java
 
b/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactory.java
index 6bf87d2..98b2fb6 100644
--- 
a/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactory.java
+++ 
b/flink-connector-mongodb/src/main/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactory.java
@@ -47,6 +47,7 @@ import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.FIL
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.LOOKUP_RETRY_INTERVAL;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_CURSOR_NO_TIMEOUT;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_FETCH_SIZE;
+import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_RECORD_SIZE;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_SAMPLES;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_SIZE;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_STRATEGY;
@@ -87,6 +88,7 @@ public class MongoDynamicTableFactory
         optionalOptions.add(SCAN_PARTITION_STRATEGY);
         optionalOptions.add(SCAN_PARTITION_SIZE);
         optionalOptions.add(SCAN_PARTITION_SAMPLES);
+        optionalOptions.add(SCAN_PARTITION_RECORD_SIZE);
         optionalOptions.add(BUFFER_FLUSH_MAX_ROWS);
         optionalOptions.add(BUFFER_FLUSH_INTERVAL);
         optionalOptions.add(DELIVERY_GUARANTEE);
diff --git 
a/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/source/MongoSourceITCase.java
 
b/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/source/MongoSourceITCase.java
index ae2d0fd..b866ad7 100644
--- 
a/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/source/MongoSourceITCase.java
+++ 
b/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/source/MongoSourceITCase.java
@@ -266,7 +266,8 @@ class MongoSourceITCase {
                 Arguments.of(PartitionStrategy.SPLIT_VECTOR, TEST_COLLECTION),
                 Arguments.of(PartitionStrategy.SAMPLE, TEST_COLLECTION),
                 Arguments.of(PartitionStrategy.SHARDED, 
TEST_SHARDED_COLLECTION),
-                Arguments.of(PartitionStrategy.SHARDED, 
TEST_HASHED_KEY_SHARDED_COLLECTION));
+                Arguments.of(PartitionStrategy.SHARDED, 
TEST_HASHED_KEY_SHARDED_COLLECTION),
+                Arguments.of(PartitionStrategy.PAGINATION, TEST_COLLECTION));
     }
 
     private static MongoSourceBuilder<RowData> defaultSourceBuilder(String 
collection) {
diff --git 
a/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactoryTest.java
 
b/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactoryTest.java
index f9084a7..3483c72 100644
--- 
a/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactoryTest.java
+++ 
b/flink-connector-mongodb/src/test/java/org/apache/flink/connector/mongodb/table/MongoDynamicTableFactoryTest.java
@@ -50,6 +50,7 @@ import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.FIL
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.LOOKUP_RETRY_INTERVAL;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_CURSOR_NO_TIMEOUT;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_FETCH_SIZE;
+import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_RECORD_SIZE;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_SAMPLES;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_SIZE;
 import static 
org.apache.flink.connector.mongodb.table.MongoConnectorOptions.SCAN_PARTITION_STRATEGY;
@@ -144,6 +145,35 @@ class MongoDynamicTableFactoryTest {
         assertThat(actual).isEqualTo(expected);
     }
 
+    @Test
+    void testMongoPaginationPartitionProperties() {
+        Map<String, String> properties = getRequiredOptions();
+        properties.put(SCAN_PARTITION_STRATEGY.key(), "pagination");
+        properties.put(SCAN_PARTITION_RECORD_SIZE.key(), "42");
+        properties.put(FILTER_HANDLING_POLICY.key(), "never");
+
+        DynamicTableSource actual = createTableSource(SCHEMA, properties);
+
+        MongoConnectionOptions connectionOptions = getConnectionOptions();
+        MongoReadOptions readOptions =
+                MongoReadOptions.builder()
+                        .setPartitionStrategy(PartitionStrategy.PAGINATION)
+                        .setPartitionRecordSize(42)
+                        .build();
+
+        MongoDynamicTableSource expected =
+                new MongoDynamicTableSource(
+                        connectionOptions,
+                        readOptions,
+                        null,
+                        LookupOptions.MAX_RETRIES.defaultValue(),
+                        LOOKUP_RETRY_INTERVAL.defaultValue().toMillis(),
+                        FilterHandlingPolicy.NEVER,
+                        SCHEMA.toPhysicalRowDataType());
+
+        assertThat(actual).isEqualTo(expected);
+    }
+
     @Test
     void testMongoLookupProperties() {
         Map<String, String> properties = getRequiredOptions();
@@ -248,6 +278,12 @@ class MongoDynamicTableFactoryTest {
                 "0ms",
                 "The 'lookup.retry.interval' must be larger than 0.");
 
+        // record size per partition should be larger than 0
+        assertSourceValidationRejects(
+                SCAN_PARTITION_RECORD_SIZE.key(),
+                "0",
+                "The record size per partition must be larger than 0.");
+
         // sink retries shouldn't be negative
         assertSinkValidationRejects(
                 SINK_MAX_RETRIES.key(),

Reply via email to