[ 
https://issues.apache.org/jira/browse/FLINK-40258?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

soyeong choe closed FLINK-40258.
--------------------------------
    Resolution: Duplicate

> MongoDB CDC SampleBucketSplitStrategy calculates negative $sample size when 
> collection count exceeds Integer.MAX_VALUE
> ----------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40258
>                 URL: https://issues.apache.org/jira/browse/FLINK-40258
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.6.0
>         Environment: * Flink CDC version: 3.6.0
> * Connector: MongoDB CDC
> * scan.incremental.snapshot.enabled: true
> * scan.incremental.snapshot.chunk.samples: 20
> * splitVector privilege: not granted
> * MongoDB collection document count: > Integer.MAX_VALUE
>            Reporter: soyeong choe
>            Priority: Major
>
> h3. Problem
> In the MongoDB CDC connector, when the splitVector command is unauthorized or 
> fails, the connector falls back to SampleBucketSplitStrategy.
> For very large MongoDB collections whose document count exceeds 
> Integer.MAX_VALUE, SampleBucketSplitStrategy calculates a negative sample 
> size due to integer overflow. This negative value is then passed to MongoDB 
> $sample, causing the snapshot split generation to fail.
> h3. Affected code
> Current master:
> [https://github.com/apache/flink-cdc/blob/master/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/assigners/splitters/SampleBucketSplitStrategy.java#L90-L97]
> Release 3.6.0:
> [https://github.com/apache/flink-cdc/blob/release-3.6.0/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mongodb-cdc/src/main/java/org/apache/flink/cdc/connectors/mongodb/source/assigners/splitters/SampleBucketSplitStrategy.java#L90-L97]
> Current code:
> {code:java}
> long count = splitContext.getDocumentCount();
> int numberOfSamples;
> if (count < DEFAULT_SAMPLING_THRESHOLD) {
>     numberOfSamples = (int) count;
> } else {
>     numberOfSamples = Math.min(numChunks * splitContext.getSamplesPerChunk(), 
> (int) count);
> }
> {code}
> The value count is a long, but it is cast to int before calculating 
> numberOfSamples. When count is greater than Integer.MAX_VALUE, this cast can 
> overflow to a negative value.
> h3. Example
> For example, a collection count around 3,727,914,087 overflows to 
> -567,053,209 when cast to int.
> The connector log shows:
> {noformat}
> Unauthorized to execute splitVector command: not authorized on <db> to 
> execute command { splitVector: "<db>.<collection>", ... }, fallback to 
> SampleSplitter
> Collection <db>.<collection> going to sample -567053209 records into 30150 
> chunks
> {noformat}
> Then MongoDB fails with:
> {noformat}
> Caused by: com.mongodb.MongoCommandException: Command failed with error 
> 28747: 'size argument to $sample must not be negative'
> {noformat}
> h3. Expected behavior
> SampleBucketSplitStrategy should calculate a positive sample size using long 
> arithmetic.
> For example, with:
> {noformat}
> numChunks = 30150
> samplesPerChunk = 20
> count = 3727914087
> {noformat}
> The expected sample size should be:
> {noformat}
> min(30150 * 20, 3727914087) = 603000
> {noformat}
> h3. Actual behavior
> The sample size becomes negative because count is cast to int before Math.min:
> {noformat}
> Math.min(603000, (int) 3727914087)
> = Math.min(603000, -567053209)
> = -567053209
> {noformat}
> This negative value is passed to MongoDB $sample, and MongoDB rejects it.
> h3. Reproduction conditions
> This issue occurs when all of the following are true:
>  * MongoDB CDC connector is used.
>  * scan.incremental.snapshot.enabled is enabled.
>  * The collection is split through SampleBucketSplitStrategy.
>  * splitVector is unauthorized or fails, so the connector falls back to 
> Sample Splitter.
>  * The MongoDB collection document count is greater than Integer.MAX_VALUE.
> h3. Suggested fix
> Use long arithmetic when calculating the requested sample size, and only 
> convert to int after validating the range.
> For example:
> {code:java}
> long requestedSamples =
>         Math.min((long) numChunks * splitContext.getSamplesPerChunk(), count);
> checkArgument(
>         requestedSamples <= Integer.MAX_VALUE,
>         "The requested sample size exceeds Integer.MAX_VALUE: %s",
>         requestedSamples);
> int numberOfSamples = (int) requestedSamples;
> {code}
> It may also be safer to calculate chunkSizeInBytes using long arithmetic:
> {code:java}
> long chunkSizeInBytes = splitContext.getChunkSizeMB() * 1024L * 1024L;
> {code}
> instead of:
> {code:java}
> long chunkSizeInBytes = splitContext.getChunkSizeMB() * 1024 * 1024;
> {code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to