This is an automated email from the ASF dual-hosted git repository.
jiabaosun pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-mongodb.git
The following commit(s) were added to refs/heads/main by this push:
new 123334b [FLINK-37255] Fix unable to configure
scan.partition.record-size option (#51)
123334b is described below
commit 123334b8a70e5b3cccfeab6b8e083fbcd7e617d2
Author: yuxiqian <[email protected]>
AuthorDate: Wed Feb 5 17:46:14 2025 +0800
[FLINK-37255] Fix unable to configure scan.partition.record-size option
(#51)
---
.../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(),