Skip to content

[fix][broker] Avoid blocking metadata read on the IO thread when redirecting migrated producers/consumers - #26051

Merged
lhotari merged 1 commit into
apache:masterfrom
merlimat:mmerli/fix-topic-migration-blocking-io
Jun 18, 2026
Merged

lhotari merged 1 commit into
apache:masterfrom
merlimat:mmerli/fix-topic-migration-blocking-io

Conversation

@merlimat

Copy link
Copy Markdown
Contributor

Motivation

AbstractTopic.getMigratedClusterUrl() resolved the migrated-cluster URL by blocking on getMigratedClusterUrlAsync().get(operationTimeoutSeconds). Its clean, non-throwing signature hid that it blocks, and several callers run on threads where blocking is harmful:

  • ServerCnx subscribe-success path, inside .thenAcceptAsync(..., ctx.executor()) (the connection's Netty IO thread), via Consumer.checkAndApplyTopicMigration();
  • the producer/consumer migration-redirect exceptionallyAsync(..., ctx.executor()) handlers in ServerCnx;
  • AbstractBaseDispatcher.checkAndApplyReachedEndOfTopicOrTopicMigration() on the dispatch path;
  • PersistentTopic / NonPersistentTopic.checkClusterMigration() and the topic write-failed path.

A slow or contended metadata read could therefore stall a Netty IO thread for up to the metadata operation timeout (default 30s). Blocking with get() on one pool while the future is completed on another is also a deadlock hazard if the completing pool is saturated.

Modifications

Route every caller through the already-existing async getMigratedClusterUrlAsync() and delete the blocking getMigratedClusterUrl() overloads:

  • Consumer.checkAndApplyTopicMigration() → checkAndApplyTopicMigrationAsync() returning CompletableFuture<Boolean> (reuses topicMigrated(...) for the send + disconnect).
  • ServerCnx: subscribe-success uses thenCompose so the migration check stays async on ctx.executor(); the three producer/consumer migration-redirect handlers resolve the URL asynchronously and then redirect-or-fail in small helpers that preserve each branch's exact cleanup and logging.
  • AbstractBaseDispatcher, PersistentTopic, NonPersistentTopic: resolve the URL asynchronously. checkClusterMigration resolves it once and reuses it for each consumer, dropping a redundant per-consumer blocking re-resolve.

A failed metadata resolution is no longer swallowed into "not migrated"; it surfaces instead. No public/admin API or wire-protocol change — getMigratedClusterUrlAsync() already existed and already hops completion to pulsar.getExecutor() to avoid the metadata-thread deadlock.

Verifying this change

This change is already covered by existing tests: ClusterMigrationTest — producer & consumer redirect, persistent & non-persistent topics, replication-backlog, resource-created, and namespace-migration cases (18 executions, all passing locally).

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

…recting migrated producers/consumers

AbstractTopic.getMigratedClusterUrl() resolved the migrated cluster URL by
blocking on getMigratedClusterUrlAsync().get(timeout). It was reached from
Netty event-loop threads and CompletableFuture callbacks (ServerCnx
subscribe/produce handlers via Consumer.checkAndApplyTopicMigration(),
AbstractBaseDispatcher, and the persistent/non-persistent checkClusterMigration
paths), so a slow or contended metadata read could stall the IO thread for up
to the metadata operation timeout.

Route all callers through the existing async getMigratedClusterUrlAsync() and
delete the blocking method:
- Consumer.checkAndApplyTopicMigration() -> checkAndApplyTopicMigrationAsync()
- ServerCnx subscribe-success uses thenCompose; the producer/consumer
  migration-redirect exceptionally handlers resolve the URL asynchronously
- checkClusterMigration resolves the URL once and reuses it for each consumer

A failed metadata resolution is no longer swallowed into 'not migrated'.
@merlimat
merlimat requested a review from lhotari June 18, 2026 02:04
@lhotari
lhotari merged commit ac053c3 into apache:master Jun 18, 2026
44 checks passed
@lhotari lhotari added this to the 5.0.0-M2 milestone Jun 18, 2026
lhotari pushed a commit that referenced this pull request Jun 22, 2026
…recting migrated producers/consumers (#26051)

(cherry picked from commit ac053c3)
lhotari pushed a commit that referenced this pull request Jun 22, 2026
…recting migrated producers/consumers (#26051)

(cherry picked from commit ac053c3)
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
…recting migrated producers/consumers (apache#26051)

(cherry picked from commit ac053c3)
nodece pushed a commit to ascentstream/pulsar that referenced this pull request Aug 28, 2026
…recting migrated producers/consumers (apache#26051)

(cherry picked from commit ac053c3)
Signed-off-by: Zixuan Liu <[email protected]>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants