SEPURI-SAI-KRISHNA opened a new pull request, #12396:
URL: https://github.com/apache/seatunnel/pull/12396

   ### Purpose of this pull request
   
   Follow-up to #12393, where @DanielLeens recorded this as Issue 1 (Low, 
non-blocking) and asked for it as a separate PR because it changes test 
behaviour. It is his preferred Option A.
   
   `testSourceRocketMqTextTagToConsole` and 
`testSourceRocketMqTextErrorTagToConsole` used fixed topic names and did not 
set a consumer group, so they fell back to the connector's default 
`SeaTunnel-Consumer-Group` (`RocketMqSourceOptions.java:30,82`) with the 
default `CONSUME_FROM_GROUP_OFFSETS` (`:53`). The broker container is created 
in `@BeforeAll` and shared across every engine leg of the template, so all legs 
of a test shared one topic and one committed offset while asserting an exact 
row count.
   
   That only holds while every leg writes exactly 32 messages and advances the 
shared offset by exactly 32. If a leg fails after writing but before 
committing, the next leg sees more than 32 unconsumed messages and fails its 
exact-count assertion, so a single bad leg can cascade into the ones after it. 
That is plausibly a contributor to the historical failure rate of these two 
tests recorded in #12322, though I have not demonstrated that link.
   
   ### The change
   
   Both tests now derive a unique topic and consumer group per invocation and 
pass them as job variables:
   
   ```java
   final String uniqueSuffix = uniqueTestSuffix();
   final String topic = "test_topic_text_tag_" + uniqueSuffix;
   final String consumerGroup = "SeaTunnel-Consumer-Group-" + uniqueSuffix;
   ...
   container.executeJob(
           "/rocketmq-source_text_tag_to_console.conf",
           Arrays.asList("sourceTopic=" + topic, "consumerGroup=" + 
consumerGroup));
   ```
   
   Topic and group share one suffix, following `testSourceRocketMqRestore:670`, 
so a failing leg's topic and group are correlatable in the logs.
   
   and the two confs read them with defaults:
   
   ```conf
   topics = "${sourceTopic:test_topic_text_tag}"
   consumer.group = "${consumerGroup:SeaTunnel-Consumer-Group}"
   ```
   
   This is not a new mechanism. It is the pattern 
`testSourceRocketMqStartConfig` and `executeRocketMqGroupOffsetsToConsole` 
already use in this class, whose Javadoc states the same purpose: "Uses 
isolated topic and consumer-group names so template invocations cannot share 
offsets." That reference test carries no `@DisabledOnContainer`, so the 
job-variable path is already exercised on every engine leg these two tests run 
on, and `executeJob(String, List<String>)` is implemented by the Zeta, Spark 
and Flink container bases alike.
   
   One consequence worth naming: a fresh consumer group has no committed 
offsets, so both tests now take the `CONSUME_FROM_GROUP_OFFSETS` cold-start 
branch and read from the first offset every time, instead of depending on the 
shared group's offset having advanced by exactly 32 since the previous leg. 
They read the 32 messages they just wrote. That is the point of the change, and 
it also means the result no longer varies with what earlier legs did.
   
   The placeholder defaults are the previous literal topic names and the 
connector's own default group, so running either conf on its own, without job 
variables, behaves exactly as it does today. The `consumer.group` line is new 
in these two confs, but its default is the value the connector already 
resolved, so it changes nothing by itself.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. Test-only: one E2E class and two test resources. No production code, no 
option, no documented behaviour.
   
   ### How was this patch tested?
   
   `./mvnw -q -DskipTests verify -pl 
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e` on JDK 11 
passes, covering the enforcer checks, `spotless:check` and compilation, and 
`test-compile` passes. `spotless:check` was also run separately with up-to-date 
caching disabled.
   
   The meaningful signal is this PR's own `rocketmq-connector-it` legs, since 
the change is about cross-leg state that only appears against a live broker. 
What should be visible is that each leg now creates its own topic, which also 
means neither test reuses data from an earlier leg.
   
   ### Relationship to #12393
   
   #12393 removes `deleteTopicIfExist` and its two call sites from the same two 
test methods. The two changes do not overlap in intent, but they sit a few 
lines apart, so whichever merges second will need a small rebase. I am happy to 
do that either way round; #12393 is the older PR and is already green, so I 
would expect this one to rebase onto it.
   


-- 
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