Skip to content

Guard streaming metrics accumulator registration - #181

Merged
cirquare merged 1 commit into
li-release-2.75from
yf186035/beam-2.75-phase4a-accumulator-guard
Oct 6, 2026
Merged

cirquare merged 1 commit into
li-release-2.75from
yf186035/beam-2.75-phase4a-accumulator-guard

Conversation

@cirquare

@cirquare cirquare commented Oct 6, 2026

Copy link
Copy Markdown

Summary

Replay/adapt the LI streaming accumulator lifecycle fix for Beam 2.75. In streaming Flink DoFnOperator, task cancellation/failover close no longer registers the full Beam metrics accumulator; metrics are registered for pipeline results only after finish()/end-of-input. Batch behavior is preserved.

This prevents cancelled streaming subtasks from attaching large Beam metrics state to Flink final execution state, which can make JobManager status snapshots merge/render huge accumulator payloads during failover recovery.

Testing Done

  • ./gradlew :runners:flink:1.20:test --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationOnOperatorClose --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationOnCancelInBatch --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationAfterFinishInStreaming --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testNoAccumulatorRegistrationOnCancelInStreaming --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:2.2:test --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationOnOperatorClose --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationOnCancelInBatch --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testAccumulatorRegistrationAfterFinishInStreaming --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest.testNoAccumulatorRegistrationOnCancelInStreaming --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:1.20:test :runners:flink:2.2:test --tests org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:1.19:build --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:1.20:build :runners:flink:2.2:build --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:2.0:build --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • ./gradlew :runners:flink:2.1:build --no-daemon -Dorg.gradle.java.home=$(/usr/libexec/java_home -v 11.0.21)
  • git diff --check

Note: an initial optional combined :runners:flink:1.19:build :runners:flink:2.1:build run timed out once in PortableStateExecutionTest > testExecution[streaming: false]; rerunning the failing test class and then full :runners:flink:1.19:build passed.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
@cirquare
cirquare marked this pull request as ready for review October 6, 2026 17:13
@cirquare
cirquare merged commit 6acd56b into li-release-2.75 Oct 6, 2026
1 check passed
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.

1 participant