rucciva opened a new pull request, #11460:
URL: https://github.com/apache/seatunnel/pull/11460
[Feature] [Connector-V2] Add NATS JetStream sink connector
### Purpose of this pull request
This pull request implements the sink portion of #10101 by adding a
publish-only NATS JetStream Connector V2 sink.
The connector supports:
- JSON and native message formats
- NATS username/password and token authentication
- Configurable subjects
- Configurable native field mappings for message ID, subject, headers, and
payload
- Optional SeaTunnel row-kind headers
- Configuration and schema validation
- At-least-once writer semantics without checkpoint state or transactional
commit support
- Pre-provisioned JetStream streams
The NATS source connector and stream-management functionality are
intentionally out of scope.
### Does this PR introduce _any_ user-facing change?
Yes. This PR adds a new `NatsJetStream` sink connector and its documentation
in English and Chinese.
Example configuration:
```hocon
sink {
NatsJetStream {
url = "nats://localhost:4222"
subject = "events"
format = "json"
}
}
```
The connector also supports native messages with configurable field mappings:
```hocon
sink {
NatsJetStream {
url = "nats://localhost:4222"
format = "native"
subject = "events"
native_format_fields {
id = "message_id"
subject = "event_subject"
headers = "attributes"
data = "payload"
}
}
}
```
This is a new connector and does not change the behavior of existing
connectors or configurations.
### How was this patch tested?
Added unit tests covering:
- Sink factory creation and option rules
- Authentication validation
- JSON and native format validation
- Native field mapping validation
- Subject fallback and resolution
- Message ID and header serialization
- Payload type validation
- Row-kind header behavior
- Invalid record handling
Executed successfully with Java 11:
```bash
./mvnw -B \
-pl seatunnel-connectors-v2/connector-nats-jetstream \
-Dtest='*NatsJetStream*Test' test
```
Result: 31 tests passed, 0 failures, 0 errors.
Executed the NATS JetStream E2E test suite with Testcontainers and NATS
`2.10.14-alpine`:
```bash
./mvnw -B \
-pl seatunnel-e2e/seatunnel-connector-v2-e2e/connector-nats-jetstream-e2e \
-am \
-DskipUT \
-DskipIT=false \
-Dit.test=NatsJetStreamIT \
-DfailIfNoTests=false \
verify
```
Result: 21 tests passed, 0 failures, 0 errors.
Also verified:
- Java 11 Spotless formatting
- Connector module compilation and verification
- Distribution packaging
- Connector registration checks
- CI module detection for the connector and E2E module
- No whitespace errors
The repository-wide license check was attempted but could not complete
because of the unrelated unresolved dependency
`connector-http-airtable:3.0.0-SNAPSHOT`.
### Check list
* [x] If any new Jar binary package adding in your PR, please add License
Notice according to the [New License
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
* [x] If necessary, please update the documentation to describe the new
feature. https://github.com/apache/seatunnel/tree/dev/docs
* [x] If necessary, please update `incompatible-changes.md` to describe the
incompatibility caused by this PR.
* [x] If you are contributing the connector code, please check that the
following files are updated:
1. Updated `plugin-mapping.properties` with the NATS JetStream connector.
2. Updated `seatunnel-dist/pom.xml`.
3. Added the NATS JetStream CI label scope in `label-scope-conf.yml`.
4. Added the NATS JetStream E2E module and test cases.
5. Updated `config/plugin_config`.
--
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]