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]