goutamadwant opened a new issue, #12166:
URL: https://github.com/apache/seatunnel/issues/12166
### Search before asking
- [x] I searched existing issues and PRs and found no dedicated Kafka source
connectivity dry-run implementation. The related framework work is listed below.
### Description
Add Kafka source support to the existing `--dry-run connect` mode so users
can check broker authentication and topic metadata before submitting a job.
On `dev` baseline `7a58d1005b864ce787ded4f0d64c6aab1135393b`,
`KafkaSourceFactory` does not implement `SupportSourceDryRunValidation`. The
connect validator therefore cannot perform a real Kafka connectivity check and
reports it as skipped. Static validation or a supplied schema does not
establish broker access or topic existence. The same gap remains at
`72acda57042efa3408681dc2c93663910b06b3e8`.
This is connector adoption of the existing opt-in SPI, not a proposal to
change the core dry-run contract.
#### Reproduction
At the baseline above, add this assertion to the existing `KafkaFactoryTest`:
```java
@Test
void testSourceSupportsConnectDryRun() {
Assertions.assertTrue(
new KafkaSourceFactory()
instanceof
org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation);
}
```
Run the focused test with the reactor prerequisites:
```bash
mvn -pl seatunnel-connectors-v2/connector-kafka -am \
'-Dtest=KafkaFactoryTest#testSourceSupportsConnectDryRun' \
-Dsurefire.failIfNoSpecifiedTests=false \
-Dskip.spotless=true -Dcheckstyle.skip verify
```
Observed on Java 11.0.19: the assertion fails with `expected: <true> but
was: <false>`, with no test error. The unchanged factory was also compiled and
checked independently on Java 8u172; the same missing-SPI assertion failed. The
Java 8 reproduction was an isolated factory assertion, not a complete reactor
test run.
This demonstrates the factory contract used by the connect validator. It is
not a claim that a complete Zeta CLI process or data job was run for the
baseline reproduction.
#### Proposed behavior and boundaries
- Reuse `KafkaSourceConfig` for runtime-equivalent output schemas, formats,
headers and event-time metadata.
- Use AdminClient topic metadata with the configured bootstrap and security
properties. Describe explicit topics; list and resolve topic patterns using
runtime matching. Support both `tables_configs` and the legacy `table_list`.
- Propagate authentication, authorization, missing-topic and timeout
failures. Preserve interruption and close the client on success and failure.
- Share a metadata-request budget of at most 30 seconds, honoring smaller
valid timeout settings and Kafka's existing timeout validation. Client setup
can take additional time.
- Do not create consumers or producers, read/write records, access or commit
consumer offsets, create consumer groups, or create topics.
- Keep unmatched patterns valid for sources waiting for future topics. In
that case, successful validation proves listing access only.
- State explicitly that metadata access does not prove consumer/group
permissions or successful message deserialization. Kafka sink support is out of
scope.
- Keep normal source execution, existing defaults, engine/API contracts and
dependencies unchanged.
#### Prepared implementation validation
- Kafka connector suite: 98 tests passed on each of Java 8u172 and Java
11.0.19.
- Authenticated Kafka 7.0.9 integration: 3 tests passed on each JDK,
covering literal/pattern/no-match success, invalid credentials, missing topics
and unchanged topic/group/record-offset probes.
- Factory discovery is covered by the Kafka unit suite. Broker tests
exercise the factory metadata contract directly, not a full Zeta job.
- Scoped Spotless and whitespace checks passed.
- Repository test-naming and Markdown checks: 5 tests passed on Java 11.
- Full-repository `./mvnw -q -DskipTests verify` passed on Java 11. Tests
were explicitly skipped in this broad build; the test executions are listed
separately above.
### Usage Scenario
Operators and deployment pipelines can detect incorrect Kafka credentials or
absent explicit topics before allocating job resources and starting a
synchronization job. This extends existing preflight coverage without reading
production messages or changing consumer-group state.
### Related issues
- #10681 tracks the broader progressive dry-run design. This issue is a
Kafka-specific follow-up, not a replacement for that umbrella.
- #11186 merged the separate opt-in connectivity SPI and explicitly leaves
connector adoption to focused follow-up PRs.
### Are you willing to submit a PR?
- [x] Yes, I am willing to submit a PR. The implementation and regression
coverage are prepared.
### Code of Conduct
- [x] I agree to follow this project's [Code of
Conduct](https://www.apache.org/foundation/policies/conduct).
--
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]