|
41 | 41 | import static org.apache.flink.cdc.connectors.base.options.SourceOptions.SCAN_NEWLY_ADDED_TABLE_ENABLED; |
42 | 42 | import static org.apache.flink.cdc.connectors.base.options.SourceOptions.SCAN_STARTUP_MODE; |
43 | 43 | import static org.apache.flink.cdc.connectors.base.options.SourceOptions.SCAN_STARTUP_TIMESTAMP_MILLIS; |
44 | | -import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.*; |
| 44 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.BATCH_SIZE; |
| 45 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.COLLECTION; |
| 46 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.CONNECTION_OPTIONS; |
| 47 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.DATABASE; |
| 48 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.FULL_DOCUMENT_PRE_POST_IMAGE; |
| 49 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.HEARTBEAT_INTERVAL_MILLIS; |
| 50 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.HOSTS; |
| 51 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.INITIAL_SNAPSHOTTING_MAX_THREADS; |
| 52 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.INITIAL_SNAPSHOTTING_PIPELINE; |
| 53 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.INITIAL_SNAPSHOTTING_QUEUE_SIZE; |
| 54 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.PASSWORD; |
| 55 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.POLL_AWAIT_TIME_MILLIS; |
| 56 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.POLL_MAX_BATCH_SIZE; |
| 57 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.RECORDS_PER_SECOND; |
| 58 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.SCAN_INCREMENTAL_SNAPSHOT_CHUNK_SAMPLES; |
| 59 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.SCAN_INCREMENTAL_SNAPSHOT_CHUNK_SIZE_MB; |
| 60 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.SCAN_INCREMENTAL_SNAPSHOT_ENABLED; |
| 61 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.SCAN_NO_CURSOR_TIMEOUT; |
| 62 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.SCHEME; |
| 63 | +import static org.apache.flink.cdc.connectors.mongodb.source.config.MongoDBSourceOptions.USERNAME; |
45 | 64 | import static org.apache.flink.cdc.debezium.utils.ResolvedSchemaUtils.getPhysicalSchema; |
46 | 65 | import static org.apache.flink.util.Preconditions.checkArgument; |
47 | 66 | import static org.apache.flink.util.Preconditions.checkNotNull; |
|
0 commit comments