WXPNG opened a new pull request, #12358:
URL: https://github.com/apache/seatunnel/pull/12358
### Purpose of this pull request
This PR adds dry-run validation support for **MySQL CDC Source** and
**Elasticsearch Sink** connectors by implementing the
`SupportSourceDryRunValidation` and `SupportSinkDryRunValidation` SPIs
respectively.
Closes #12357
### Details
#### MySQL CDC Source (`MySqlIncrementalSourceFactory`)
- Implements `SupportSourceDryRunValidation` interface
- **`inferSchemaForDryRun()`**: Loads the MySQL JDBC driver and reads table
metadata to infer schema, which implicitly validates connectivity and basic
SELECT privilege
- **`validateConnectionForDryRun()`**: Creates a real MySQL connection via
Debezium's `MySqlConnection` and validates that the configured user has:
- `REPLICATION SLAVE` privilege (required for reading binlog events)
- `REPLICATION CLIENT` privilege (required for `SHOW MASTER STATUS`)
- Permission validation is also performed in `restoreSource()` so issues
surface at task submission time (both HTTP API and CLI)
- Error messages include actionable `GRANT` statements for quick remediation
#### Elasticsearch Sink (`ElasticsearchSinkFactory`)
- Implements `SupportSinkDryRunValidation` interface
- **`validateConnectionForDryRun()`**: Creates an `EsRestClient` and
performs:
1. Cluster connectivity & authentication check via `GET /` (fetches
cluster info)
2. Target index accessibility check via `HEAD /{index}`
- Skips index existence check for dynamic index names containing `${...}`
placeholders (e.g., `seatunnel_${age}`)
- Adds `INDEX_VARIABLE_PREFIX` constant to `ElasticsearchSinkOptions`
- Connection validation is also performed in `createSink()` for early
failure at task submission time
### Changes Summary
| File | Change |
|------|--------|
| `MySqlIncrementalSourceFactory.java` | Implement
`SupportSourceDryRunValidation`, add CDC permission validation |
| `MySqlIncrementalSourceFactoryTest.java` | Add unit tests for dry-run
validation |
| `ElasticsearchSinkFactory.java` | Implement `SupportSinkDryRunValidation`,
add connection validation |
| `ElasticsearchSinkOptions.java` | Add `INDEX_VARIABLE_PREFIX` constant |
| `ElasticsearchSinkDryRunValidationTest.java` | Add unit tests for dry-run
validation (new file) |
### How to test
1. **MySQL CDC Source**: Configure a MySQL CDC source job with incorrect
credentials or a user lacking REPLICATION privileges, then run with
`--dry-run`. The error should be reported immediately with an actionable GRANT
statement.
2. **Elasticsearch Sink**: Configure an Elasticsearch sink job with wrong
hosts or credentials, then run with `--dry-run`. The connection error should be
reported immediately.
### Checklist
- [x] Code formatted with `./mvnw spotless:apply`
- [x] Unit tests added and passing
- [x] No incompatible changes to existing APIs
- [x] Apache License header present on all new files
--
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]