Skip to content

[fix][broker] Handle synchronous schema lookup failures in replication - #26108

Open
Denovo1998 wants to merge 6 commits into
apache:masterfrom
Denovo1998:handle_synchronous_schema_lookup_failures_in_replication
Open

[fix][broker] Handle synchronous schema lookup failures in replication#26108
Denovo1998 wants to merge 6 commits into
apache:masterfrom
Denovo1998:handle_synchronous_schema_lookup_failures_in_replication

Conversation

@Denovo1998

@Denovo1998 Denovo1998 commented Jun 29, 2026

Copy link
Copy Markdown
Contributor

Motivation

Geo replication pauses and rewinds the cursor when a replicated message needs schema information that is not immediately available. However, if the local schema lookup throws synchronously before returning a future, the current entry is not cleaned up through the schema-fetch path.

This can leave the in-flight task permit incomplete and skip releasing the entry resources for the failed message.

Modifications

  • Catch synchronous failures from getSchemaInfo(msg) in GeoPersistentReplicator.
  • Release the current entry, retained payload buffer, and recycled message on that failure path.
  • Mark the current in-flight entry as completed before rewinding the cursor.
  • Add a regression test covering synchronous schema lookup failure cleanup.
  • Normalize checked, unchecked, and Guava-wrapped synchronous schema lookup
    failures into failed futures.
  • Back off by MESSAGE_RATE_BACKOFF_MS before rewinding after a failed schema fetch.
  • Skip the outer readMoreEntries() call while waiting for cursor rewind, leaving
    doRewindCursor(true) responsible for resuming reads.

Verifying this change

  • Make sure that the change passes the CI checks.

  • gradlew :pulsar-broker:test --tests org.apache.pulsar.broker.service.persistent.GeoPersistentReplicatorTest -PtestRetryCount=0`

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

@void-ptr974 void-ptr974 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 the fix. The cleanup path makes sense to me.

I left a few comments around exception handling, retry behavior, and test coverage.

CompletableFuture<SchemaInfo> schemaFuture;
try {
schemaFuture = getSchemaInfo(msg);
} catch (Exception e) {

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.

Would it be better to narrow this catch to the expected exception type? Since getSchemaInfo only declares ExecutionException, catching all Exceptions could accidentally turn unrelated bugs into schema retry loops. Another option might be to normalize getSchemaInfo to return a failed future and then reuse the existing schemaFuture.isCompletedExceptionally() path.

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.

Good point. I narrowed the catch to ExecutionException and normalize that synchronous schema lookup failure into a failed future, so unexpected exceptions are no longer converted into schema retry loops while the existing schema future cleanup path is reused.

headersAndPayload.release();
msg.recycle();
skipRemainingMessages = true;
doRewindCursor(false);

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.

This path can immediately rewind and re-read the same entry if the synchronous schema lookup failure persists. replicateEntries() returns false, so readEntriesComplete() may call readMoreEntries() right away.

One way to avoid a tight retry loop is to keep the replicator in the cursor-rewinding wait state and schedule doRewindCursor(true) after a small backoff instead of rewinding immediately.

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.

Agreed. The exceptional schema future path now keeps the replicator in the cursor-rewinding wait state and schedules doRewindCursor(true) after MESSAGE_RATE_BACKOFF_MS. Successful schema fetches still rewind immediately.

return null;
}).when(entry).release();

List<Entry> entries = List.of(entry);

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.

This test only covers the current entry cleanup. It does not verify the batch behavior after skipRemainingMessages is set.

Please extend it to use a multi-entry batch and verify that the remaining entries are skipped/released, completedEntries reaches the full batch size, and the cursor is rewound.

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.

Extended the regression test to use a multi-entry batch. It now verifies the remaining entry is skipped and released, completedEntries reaches the full batch size, and cursor rewind/read retry is triggered only by the scheduled backoff task.

@void-ptr974

Copy link
Copy Markdown
Contributor

Thanks for the update. The main concerns look addressed.

One small follow-up: after this path calls beforeTerminateOrCursorRewinding(...), replicateEntries() returns false, so readEntriesComplete() will still call readMoreEntries(). That call should not start a read while waitForCursorRewindingRefCnf > 0, but it may schedule another delayed retry. Could we skip the outer readMoreEntries() when the replicator is already waiting for cursor rewind, and let the scheduled doRewindCursor(true) resume reads instead?

@Denovo1998

Copy link
Copy Markdown
Contributor Author

Thanks for the update. The main concerns look addressed.

One small follow-up: after this path calls beforeTerminateOrCursorRewinding(...), replicateEntries() returns false, so readEntriesComplete() will still call readMoreEntries(). That call should not start a read while waitForCursorRewindingRefCnf > 0, but it may schedule another delayed retry. Could we skip the outer readMoreEntries() when the replicator is already waiting for cursor rewind, and let the scheduled doRewindCursor(true) resume reads instead?

I added a guard in readEntriesComplete() to skip the outer readMoreEntries() while the replicator is already waiting for cursor rewind. This leaves the scheduled doRewindCursor(true) as the path that resumes reads. I also updated the regression test to exercise readEntriesComplete() end-to-end and verify no extra read is triggered before the scheduled rewind runs.

@void-ptr974 void-ptr974 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.

LGTM. Thanks for addressing the comments.

@lhotari lhotari 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.

I ran an AI-assisted review of this PR (combined Claude Fable 5 + OpenAI Codex gpt-5.6-sol review; findings merged and verified against the code). Overall the fix looks real and correctly targeted: getSchemaInfo() is a Guava LoadingCache.get() call that throws ExecutionException synchronously, and in current master that throw lands in the outer catch (Exception) after headersAndPayload.retain() — leaking the retained buffer, the entry, the MessageImpl and the in-flight permit, with the cursor neither rewound nor resumed. Routing the failure into the existing isCompletedExceptionally() skip-path reuses the proven cleanup + rewind machinery, and the new readEntriesComplete() guard properly defers read resumption to the scheduled doRewindCursor(true).

Findings, in decreasing severity:

  1. Catching only ExecutionException leaves the same leak for unchecked synchronous failures (GeoPersistentReplicator.replicateEntries). Guava's LoadingCache.get() also throws UncheckedExecutionException (loader threw a RuntimeException) and ExecutionError, and getSchemaByVersion() can throw unchecked synchronously. Any of those still escape to the outer catch (Exception e), which only logs — reproducing exactly the failure mode this PR sets out to fix. Suggest broadening to catch (Exception e), or cleaner: move the try/catch into getSchemaInfo() so it returns a failed future and drop throws ExecutionException (a CompletableFuture-returning method shouldn't throw synchronously; GeoPersistentReplicator is its only caller and ShadowReplicator doesn't use it, so the signature change is contained).

  2. The regression test self-repairs the leak it's meant to detect (GeoPersistentReplicatorTest, the finally block). The loop releasing headersAndPayload until refCnt() == 0 means the test would still pass if the production path forgot headersAndPayload.release(). Suggest asserting assertThat(headersAndPayload.refCnt()).isZero() right after the verifications, before any fallback cleanup. If finding 1 is addressed, please also add a companion test injecting an unchecked exception (e.g. UncheckedExecutionException), which the current test cannot cover.

  3. Two behavior changes are not reflected in the PR description. The diff also (a) adds a MESSAGE_RATE_BACKOFF_MS (1s) delay before doRewindCursor(true) for all schema-fetch failures, including asynchronous ones — previously an async failure rewound immediately, so a persistently failing schema fetch could hot-loop read → fail → rewind → re-read; and (b) replaces the old 1s scheduled-retry polling from readEntriesComplete() with an explicit hand-off to the scheduled rewind. Both are good changes, but the description/commit message should state them — especially with the release/4.0.13 and release/4.2.4 labels, since backporters need the full behavioral delta (and should verify those branches have the InFlightTask/waitForCursorRewindingRefCnf structure this patch assumes).

  4. Minor: the new debug log in readEntriesComplete() reads reasonOfWaitForCursorRewinding, which is non-volatile and written under the inFlightTasks lock, so the log line can print a stale or null reason. Harmless (log-only), just confirming it's intentional.

  5. Minor: the scheduled rewind is fire-and-forget — if the broker executor rejects the task at shutdown, the exception is swallowed inside whenComplete and waitForCursorRewindingRefCnf never decrements, stalling that replicator until unload. This matches the existing idiom in readMoreEntries(), so it's acceptable as-is.

Concurrency was checked independently by both reviews and no race was found in the new guard: (a) for Fetching_Schema, resume is deterministically owned by doRewindCursor(true) (immediate on success, backoff-scheduled on failure), and if the rewind wins the race against the guard, the fallthrough readMoreEntries() safely no-ops on hasPendingRead(); (b) for the Failed_Publishing transient window (refcount briefly > 0 between beforeTerminateOrCursorRewinding and doRewindCursor(false) on the producer thread), resumption is still guaranteed because the same sendComplete continues on its thread and its queue-drain logic calls readMoreEntries() after the refcount is back to 0; (c) Terminating is handled by the earlier state check in readEntriesComplete().

@Denovo1998

Copy link
Copy Markdown
Contributor Author

I ran an AI-assisted review of this PR (combined Claude Fable 5 + OpenAI Codex gpt-5.6-sol review; findings merged and verified against the code). Overall the fix looks real and correctly targeted: getSchemaInfo() is a Guava LoadingCache.get() call that throws ExecutionException synchronously, and in current master that throw lands in the outer catch (Exception) after headersAndPayload.retain() — leaking the retained buffer, the entry, the MessageImpl and the in-flight permit, with the cursor neither rewound nor resumed. Routing the failure into the existing isCompletedExceptionally() skip-path reuses the proven cleanup + rewind machinery, and the new readEntriesComplete() guard properly defers read resumption to the scheduled doRewindCursor(true).

Findings, in decreasing severity:

  1. Catching only ExecutionException leaves the same leak for unchecked synchronous failures (GeoPersistentReplicator.replicateEntries). Guava's LoadingCache.get() also throws UncheckedExecutionException (loader threw a RuntimeException) and ExecutionError, and getSchemaByVersion() can throw unchecked synchronously. Any of those still escape to the outer catch (Exception e), which only logs — reproducing exactly the failure mode this PR sets out to fix. Suggest broadening to catch (Exception e), or cleaner: move the try/catch into getSchemaInfo() so it returns a failed future and drop throws ExecutionException (a CompletableFuture-returning method shouldn't throw synchronously; GeoPersistentReplicator is its only caller and ShadowReplicator doesn't use it, so the signature change is contained).
  2. The regression test self-repairs the leak it's meant to detect (GeoPersistentReplicatorTest, the finally block). The loop releasing headersAndPayload until refCnt() == 0 means the test would still pass if the production path forgot headersAndPayload.release(). Suggest asserting assertThat(headersAndPayload.refCnt()).isZero() right after the verifications, before any fallback cleanup. If finding 1 is addressed, please also add a companion test injecting an unchecked exception (e.g. UncheckedExecutionException), which the current test cannot cover.
  3. Two behavior changes are not reflected in the PR description. The diff also (a) adds a MESSAGE_RATE_BACKOFF_MS (1s) delay before doRewindCursor(true) for all schema-fetch failures, including asynchronous ones — previously an async failure rewound immediately, so a persistently failing schema fetch could hot-loop read → fail → rewind → re-read; and (b) replaces the old 1s scheduled-retry polling from readEntriesComplete() with an explicit hand-off to the scheduled rewind. Both are good changes, but the description/commit message should state them — especially with the release/4.0.13 and release/4.2.4 labels, since backporters need the full behavioral delta (and should verify those branches have the InFlightTask/waitForCursorRewindingRefCnf structure this patch assumes).
  4. Minor: the new debug log in readEntriesComplete() reads reasonOfWaitForCursorRewinding, which is non-volatile and written under the inFlightTasks lock, so the log line can print a stale or null reason. Harmless (log-only), just confirming it's intentional.
  5. Minor: the scheduled rewind is fire-and-forget — if the broker executor rejects the task at shutdown, the exception is swallowed inside whenComplete and waitForCursorRewindingRefCnf never decrements, stalling that replicator until unload. This matches the existing idiom in readMoreEntries(), so it's acceptable as-is.

Concurrency was checked independently by both reviews and no race was found in the new guard: (a) for Fetching_Schema, resume is deterministically owned by doRewindCursor(true) (immediate on success, backoff-scheduled on failure), and if the rewind wins the race against the guard, the fallthrough readMoreEntries() safely no-ops on hasPendingRead(); (b) for the Failed_Publishing transient window (refcount briefly > 0 between beforeTerminateOrCursorRewinding and doRewindCursor(false) on the producer thread), resumption is still guaranteed because the same sendComplete continues on its thread and its queue-drain logic calls readMoreEntries() after the refcount is back to 0; (c) Terminating is handled by the earlier state check in readEntriesComplete().

@lhotari
Okay, I've reviewed it again and made some changes. Take a look once more.

@void-ptr974 void-ptr974 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.

LGTM

lhotari
lhotari previously approved these changes Jul 26, 2026
@lhotari
lhotari dismissed their stale review July 26, 2026 10:54

Dropping approval while checking for possible race conditions.

Comment on lines +459 to +462
} else if (waitForCursorRewindingRefCnf > 0) {
log.debug()
.attr("reason", reasonOfWaitForCursorRewinding)
.log("Skipping read while waiting for cursor rewind");

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.

readMoreEntries already handles this case. it's better to leave it there. waitForCursorRewindingRefCnf is designed to be referenced inside a synchronized(inflightTasks) { block.

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.

@lhotari
That makes sense. Since readMoreEntries() already checks waitForCursorRewindingRefCnf under the inFlightTasks lock, the outer guard is redundant and reads the rewind state outside its intended synchronization boundary.

I will remove the guard from readEntriesComplete() and update the regression test to verify that no cursor read is scheduled while waiting for a rewind, rather than asserting that readMoreEntries() is not invoked.

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.

3 participants