danmorcov88 opened a new pull request, #11713: URL: https://github.com/apache/nifi/pull/11713
# Summary [NIFI-7180](https://issues.apache.org/jira/browse/NIFI-7180) Adds `CaptureChangeMongoDB`, a source processor that reads a MongoDB change stream and writes the events as records. NiFi has CDC for MySQL but nothing for MongoDB, and the ticket has been open since 2020 without a pull request. The usual answer today is Debezium or the MongoDB Kafka Connector plus Kafka between MongoDB and NiFi, which is two extra systems for what one change stream cursor does on its own. The bundle is laid out as `nifi-extension-bundles/nifi-cdc/nifi-cdc-mongodb-bundle`, next to the MySQL one, and uses the existing `MongoDBClientService`, so connection settings and credentials stay in one place. ## What it does - Reads insert, update, replace, delete and invalidate events for one collection or for every collection of a database, and writes them through a Record Writer, one FlowFile per batch. - Every record has the same fields, whatever the collection holds: `operation`, `database`, `collection`, `document_key`, `full_document`, `full_document_before_change`, `updated_fields`, `removed_fields`, `cluster_time`, `wall_time`, `txn_number`, `resume_token`. - An aggregation pipeline can be given and is applied by the server before the events are sent, so the events it drops never travel. Only the stages MongoDB allows on a change stream are accepted, and the pipeline is checked while the processor is configured. - `Start Position` chooses the current moment, a point in time, or an initial snapshot that writes the documents a collection already holds and then carries on from the moment the snapshot was taken. - `Full Document`, `Full Document Before Change` and `Extended JSON Mode` map onto the driver options. - Implements `VerifiableProcessor`: connection, replica set or `mongos`, server version, the `changeStream` and `find` privileges on the watched scope, and the pre-image configuration when it is asked for. ## Notes on the design **Documents are Extended JSON strings, not mapped record fields.** A collection has no fixed shape, so a schema per collection would either be wrong or would have to be inferred. Writing the document as Extended JSON keeps every BSON type and lets a flow pick it apart with the record path processors when it wants to. `Extended JSON Mode` chooses between the relaxed and the canonical form. **The resume token is stored with `session.setState` inside the session, not after the commit.** The token and the FlowFiles of its batch are then part of the same transaction, so the stored position can never run ahead of what was delivered. A failure replays the events after the last committed token and never steps over them, which makes delivery at-least-once. Before the first commit there is no committed token, so the position the cursor reports when it is opened is kept in memory and used instead; without that, a batch that failed before the first commit would reopen the stream at "now" and skip exactly the events that failed. **The stream is always opened with `startAfter`, never `resumeAfter`.** The two behave the same for an ordinary token, and only `startAfter` accepts the token of an invalidate event. That is what lets the stream carry on when the watched collection is dropped or renamed, without a second code path or a state flag recording which of the two to use. It needs MongoDB 4.2 or later; the processor requires 6.0. **A position the server cannot serve is recognised in both forms.** `ChangeStreamHistoryLost` (286) is the documented one, but a server asked to resume from a token that is not in the oplog answers `ChangeStreamFatalError` (280) with "cannot resume stream; the resume token was not found". Both count; other fatal stream errors keep the ordinary retry. `On History Lost` decides whether the flow stops with an error and keeps the position, which is the default, or gives it up on purpose and restarts at the present. **An idle stream still stores progress.** The server reports how far it has read even when it sends no event, so a quiet collection does not keep an old position while the oplog moves on. **Deployment scope is not included.** A change stream over a whole deployment needs `MongoClient.watch()`, which the driver sends with `allChangesForCluster: true`. `MongoDBClientService` hands out a `MongoDatabase` only, and `getDatabase("admin").watch()` is not the same thing, since the driver tags it at database level and it reports changes to `admin` alone. Collection and database scope are covered. If the project is open to adding an accessor for the `MongoClient` to `MongoDBClientService`, the field already exists on `MongoDBControllerService` and deployment scope becomes a third branch; I am happy to raise that as a separate Jira rather than widen this pull request. **The processor never writes to the source.** It needs `changeStream` and `find` on the watched scope and nothing else: no collection is created, no index is added, no `collMod` is run. Pre-images are turned on by an administrator. An integration test reads `serverStatus().opcounters` before and after a run and fails if the insert, update or delete counters moved. **Dependencies.** `mongodb-driver-sync` and `nifi-mongodb-client-service-api` are `provided`; the NAR declares `nifi-mongodb-client-service-api-nar` as its parent NAR, the same way `nifi-mongodb-nar` does, so neither the driver nor the service API is bundled a second time. No new third-party dependency is introduced. ## Testing Unit tests drive the processor with a fake cursor and a clock the test moves, covering the record mapping for every operation type, the batch loop, the stored state, the pipeline validation, the behaviour after a failed batch, the two forms of a lost position, and the growing wait between attempts. Integration tests use Testcontainers against MongoDB 7.0 and 8.0 and cover: every operation in order; changes made while the processor is stopped; a dropped collection and the stream carrying on past the invalidate; a position the server rejects, with both settings of `On History Lost`; a paused server and the recovery afterwards; a multi-document transaction sharing its `txn_number`; database scope; a `$match` pipeline filtering on the server; `Update Lookup`; pre-images; canonical Extended JSON; a start position in the past; that no write reaches the server; and an initial snapshot of fifty thousand documents with writes made while it runs, where every document present at the start has to come out, plus a snapshot interrupted after three batches that carries on and reads every document exactly once. # Tracking ### Issue Tracking - [Apache NiFi Jira](https://issues.apache.org/jira/browse/NIFI-7180) issue created ### Pull Request Tracking - Pull Request title starts with Apache NiFi Jira issue number `NIFI-7180` - Pull Request commit message starts with Apache NiFi Jira issue number `NIFI-7180` - The commit is not signed with a registered key ### Pull Request Formatting - Pull Request based on the current revision of the `main` branch - Pull Request refers to a feature branch with one commit containing the changes # Verification Please indicate the verification steps performed prior to pull request creation. ### Build - [x] Build completed using `./mvnw clean install -P contrib-check` - [x] JDK 21 - [x] JDK 25 Run for the new bundle and everything it depends on, `-pl nifi-extension-bundles/nifi-cdc/nifi-cdc-mongodb-bundle/nifi-cdc-mongodb-nar -am`, 105 modules, on both JDKs. ### Licensing - [x] New dependencies are compatible with the [Apache License 2.0](https://apache.org/licenses/LICENSE-2.0) according to the [License Policy](https://www.apache.org/legal/resolved.html) - [x] New dependencies are documented in applicable `LICENSE` and `NOTICE` files No third-party dependency is added to the project. `mongodb-driver-sync` is already used by `nifi-mongodb-bundle` and is taken here with scope `provided` at the version that bundle declares; it reaches the processor at runtime through `nifi-mongodb-client-service-api-nar`, which the NAR declares as its parent NAR. Nothing new is bundled, so no `LICENSE` or `NOTICE` entry is needed; the NAR carries a `NOTICE` like `nifi-cdc-mysql-nar`. ### Documentation - [x] Documentation formatting appears as expected in rendered files The processor documentation comes from the annotations. It was checked in a running NiFi 2.x: the processor appears in the palette, and the property list, the allowable values, the dependent properties and the state description render as intended. -- 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]
