Skip to content

MINOR: Propagate ClusterInstance#createTopic errors and clean up callers - #23699

Open
lytt925 wants to merge 8 commits into
apache:trunkfrom
lytt925:followup-clusterinstance-createtopic
Open

lytt925 wants to merge 8 commits into
apache:trunkfrom
lytt925:followup-clusterinstance-createtopic

Conversation

@lytt925

@lytt925 lytt925 commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Follow-up to
#23513
to address the review comments.

ClusterInstance#createTopic and #createTopicWithAssignment did not wait
for the CreateTopics result. If the creation failed, the real error was
hidden by the 60 second waitTopicCreation timeout, or ignored if the
topic already existed. They now call get() so that the real error can be
propagated.

These are setup helpers that are expected to succeed, so
ExecutionException and InterruptedException are caught and callers no
longer need to declare them.

Other cleanups:

  • Remove assertDoesNotThrow wrappers and throws clauses that only
    existed
    because of the checked exceptions.
  • Call ClusterInstance#createTopic directly instead of pass-through
    private helpers.
  • Inline single-use local variables in TopicCommandTest and drop a
    redundant Map.of() in RemoteTopicCrudTest.
  • PlaintextConsumerTest#testSeek created the same topic twice; the
    duplicate call is removed.

Co-Authored by Claude

Reviewers: Chia-Ping Tsai chia7712@gmail.com, Ken Huang
s7133700@gmail.com

ClusterInstance#createTopic and #createTopicWithAssignment did not wait for
the CreateTopics result, so a failed creation was hidden behind the 60 second
waitTopicCreation timeout. Call get() so that the real error is surfaced.

These are setup helpers that are expected to succeed, so ExecutionException
and InterruptedException are wrapped in RuntimeException (restoring the
interrupt flag) instead of being declared.

- PlaintextConsumerTest#testSeek created the same topic twice. The second
  call used to fail silently and now throws TopicExistsException.
- Remove the assertDoesNotThrow wrappers that only existed to handle the
  checked InterruptedException.
ClusterInstance#createTopic and #createTopicWithAssignment no longer declare
checked exceptions. Remove the InterruptedException and ExecutionException
declarations, and the resulting unused imports, that only existed because of
them.
…pers

These private helpers only forwarded to ClusterInstance#createTopic or
#createTopicWithAssignment, so call the ClusterInstance methods directly.
- TopicCommandTest: inline the partition count and replication factor
  locals. They were also passed to waitTopicCreation before, but are now
  only used once in createTopic.
- RemoteTopicCrudTest: drop the redundant Map.of() argument.
@github-actions github-actions Bot added triage PRs from the community core Kafka Broker tools tests Test fixes (including flaky tests) storage Pull requests that target the storage module tiered-storage Related to the Tiered Storage feature clients labels Oct 4, 2026

@chia7712 chia7712 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@lytt925 thanks for this patch!

} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Failed to create topic " + topicName, e);
} catch (ExecutionException e) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We can throw the cause directly if it is a subtype of RuntimeException

} catch (ExecutionException e) {
    if (e.getCause() instanceof RuntimeException re) throw re;
    throw new RuntimeException("Failed to create topic " + topicName, e.getCause());
}

@lytt925 lytt925 Oct 5, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for pointing this out! Updated

} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Failed to create topic " + topicName, e);
} catch (ExecutionException e) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ditto

}

default void createTopic(String topicName, int partitions, short replicas, Map<String, String> props) throws InterruptedException {
default void createTopic(String topicName, int partitions, short replicas, Map<String, String> props) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

BTW, we could apply the same style to deleteTopic to get rid of the checked exceptions :)

@m1a2st m1a2st left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for @lytt925 update, one minor comment

Comment on lines 501 to 503
private void createTopic(String name, short replicationFactor) {
clusterInstance.createTopic(name, 1, replicationFactor);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We can inline this method.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

updated in acbe9cc

@m1a2st m1a2st left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you also remove the throws clauses that are no longer reachable? A few were missed:

  • ProducerFailureHandlingTest: createInternalTopic and its two callers
  • ResetIntegrationTest
  • TransactionsTestHelper: the five timeout / testEmptyAbortAfterCommit helper methods. The throws Exception on their call sites in TransactionsTest and TransactionsWithTieredStoreTest is redundant as well.
  • ConfigCommandIntegrationTest, ResetConsumerGroupOffsetTest, and TopicCommandTest

@github-actions github-actions Bot removed the triage PRs from the community label Oct 6, 2026
@lytt925

lytt925 commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

@m1a2st Thanks! Removed the ones you listed. I also did a broader sweep with the same rule and found a few more in AddPartitionsTest, AdminMetadataTest, CreateTopicsRequestWithPolicyTest, RemoteTopicCrudTest, and PlaintextConsumerTest (including the call sites of its private helpers).

see: 29f28a7 and 0d06245

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

ci-approved clients core Kafka Broker storage Pull requests that target the storage module tests Test fixes (including flaky tests) tiered-storage Related to the Tiered Storage feature tools

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants