Skip to content

[Improve] Consume any ready child in CollectOperator - #18813

Open
JackieTien97 wants to merge 1 commit into
apache:masterfrom
JackieTien97:codex/collect-any-child-ready
Open

JackieTien97 wants to merge 1 commit into
apache:masterfrom
JackieTien97:codex/collect-any-child-ready

Conversation

@JackieTien97

Copy link
Copy Markdown
Contributor

Description

CollectOperator currently stalls on its current child even when another child can produce data. Consume whichever child becomes unblocked, retaining the selected index on the driver thread so MappingCollectOperator applies the correct column mapping. This follows CollectNode's existing allowance to consume other children while one is blocked.

Cache each child's blocking future and register at most one listener per pending child future. Reuse the aggregate future within a wait round, and retarget existing listeners between rounds with a completion recheck to prevent missed wakeups. Already-ready children need no listener registration or aggregate-future allocation. Child failures and cancellations are propagated before consumption.

Track remaining children independently of their indexes, close all remaining children, and account for intermediate state retained by other children when estimating memory.

Validation

Both clean reactor builds passed, including dependencies and all 15 new unit tests:

mvn clean test -pl iotdb-core/calc-commons -am \
  -Dtest=AnyChildBlockedTest,CollectOperatorTest \
  -Dsurefire.failIfNoSpecifiedTests=false -DfailIfNoTests=false

mvn clean test -pl iotdb-core/calc-commons -am -P with-zh-locale \
  -Dtest=AnyChildBlockedTest,CollectOperatorTest \
  -Dsurefire.failIfNoSpecifiedTests=false -DfailIfNoTests=false

Coverage includes 2,000 wait rounds with one long-blocked child, 1,000 concurrent-completion rounds, retargeting races and delayed callbacks, cancellation/failure handling, out-of-order completion and cleanup, column mapping, and retained-memory estimates. Spotless formatting and git diff --check also passed.

  • Self-reviewed, including completion-thread/driver-thread interactions.
  • Added unit tests for new behavior and concurrency cases.
  • Added Javadocs and comments explaining listener reuse and wakeup races.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant