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]

Reply via email to