mxtymoshyk opened a new pull request, #40393:
URL: https://github.com/apache/beam/pull/40393
Fixes #18723.
Fixes a race in two `MqttIOTest` read tests and re-enables
`testReadNoClientId`, which has been `@Ignore`d under #18723. Test-only change.
**Root cause**
Before publishing, both tests called `doConnect(connection ->
!connection.getConnectionId().isEmpty())`, which waits until the broker has any
client with a non-empty connection id. The test's own publish client connects
before the publisher thread starts, so the condition was already true and the
wait returned after one 1-second poll. If the pipeline's reader had not
subscribed by then, the broker had no subscriber for the topic and dropped the
messages.
- `testReadWithMetadata` reads for at most 10 seconds and then fails with
every record missing. That is the failure on a recent `Java_IOs_Direct` run:
```
java.lang.AssertionError:
MqttIO.Read/Read(UnboundedMqttSource)/StripIds/ParMultiDo(StripIds).output:
Expected: iterable with items [<MqttRecord{topic=topic/1, ...}>, ...]
but: no item matches: <MqttRecord{topic=topic/1, ...}>, ...
```
#34133 added a 2-second sleep for this (#34175). The sleep made the race
less likely without removing it.
- `testReadNoClientId` has no max read time, so the reader waits forever for
records that never arrive and JUnit kills it at 60 seconds. That matches the
timeout reported in #18723.
**Fix**
A new helper, `waitForSubscription`, polls the broker's topic subscriptions
every 100 ms until one matches the destination the test is about to publish to.
`Subscription.matches` handles wildcards, so the `topic/#` filter in
`testReadWithMetadata` matches `topic.1`. The publish clients never subscribe,
so a match means the reader has subscribed. The helper gives up after 30
seconds.
Both tests use it in place of `doConnect`, and the 2-second sleep is gone.
The publisher threads now log exceptions instead of swallowing them, so a
failed wait shows up in the test log. `testRead` already waits for the reader's
own client id and is unchanged.
**Repro**
Add a delay to the reader in `MqttIO.UnboundedMqttReader.start()`, between
`client.connect()` and the subscribe call, to stand in for a slow runner:
```java
client.connect();
Thread.sleep(5000);
```
| | Before | After |
|---|---|---|
| `testReadWithMetadata`, reader delayed 5s | fails with `no item matches` |
passes |
| `testReadNoClientId`, reader delayed 5s | `test timed out after 60000
milliseconds` | passes |
| `testReadWithMetadata`, no delay, 20 runs | | 20/20, about 1.4s each (4.2s
before) |
| `testReadNoClientId`, no delay, 20 runs | | 20/20, about 1.4s each |
Locally, `spotlessCheck`, `checkstyleTest` and `:sdks:java:io:mqtt:test`
pass: 21 tests, none skipped.
------------------------
Thank you for your contribution! Follow this checklist to help us
incorporate your contribution quickly and easily:
- [x] Mention the appropriate issue in your description (for example:
`addresses #123`), if applicable. This will automatically add a link to the
pull request in the issue. If you would like the issue to automatically close
on merging the pull request, comment `fixes #<ISSUE NUMBER>` instead.
- [ ] Update `CHANGES.md` with noteworthy changes.
- [ ] If this contribution is large, please file an Apache [Individual
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).
See the [Contributor Guide](https://beam.apache.org/contribute) for more
tips on [how to make review process
smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier).
--
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]