Skip to content

Fix table UDF result block splitting - #18333

Open
ColinLeeo wants to merge 9 commits into
masterfrom
fix/table-udf-result-chunking
Open

Fix table UDF result block splitting#18333
ColinLeeo wants to merge 9 commits into
masterfrom
fix/table-udf-result-chunking

Conversation

@ColinLeeo

@ColinLeeo ColinLeeo commented Jul 28, 2026

Copy link
Copy Markdown
Collaborator

Description

Problem

A table-model UDF with a table argument can accumulate a large result while processing a partition, especially when one device contains many rows or the UDF expands each input row into multiple output rows. Previously, TableFunctionOperator returned each generated TsBlock directly. A sufficiently large result could therefore reach the exchange or RPC layer as one block and exceed its frame-size limit.

Pass-through columns are appended only after the UDF has produced its proper columns, so the final block must be split after that step to keep all output columns aligned.

Design

TableFunctionOperator now extends AbstractOperator and reuses its existing lazy result-splitting mechanism.

The final block is first built with any pass-through columns appended and then queued. Before a queued block is returned, checkTsBlockSizeAndGetResult() derives a maximum position count from the logical average row size of the first non-empty result block. If the block exceeds that position count, AbstractOperator retains it and returns ordered regions over subsequent next() calls. These regions are views over the original columns and do not rebuild, copy, or serialize values.

This design avoids both expensive alternatives:

  • serializing candidate regions only to serialize the selected regions again downstream;
  • copying batched UDF output row by row into another TsBlockBuilder while tracking its size.

The configured max_tsblock_size_in_bytes is therefore used as a low-cost logical target rather than an exact serialized-size guarantee. The mechanism does not inspect every variable-width payload, and one indivisible oversized row or highly skewed binary values may still exceed the target. This PR is intended to prevent large, ordinary partition results from being returned as a single block while keeping the UDF execution path inexpensive.

No new configuration is introduced. This change deliberately does not add a separate max_tsblock_line_number constraint; it follows the existing AbstractOperator behavior.

Operator lifecycle

Pending result blocks are drained before another partition is processed, and hasNext()/isFinished() continue to report pending retained or queued results. Closing the operator clears the partition cache, queued results, and retained result references before releasing the child operator and UDF resources.

Tests

The unit tests cover:

  • oversized output produced by process(), with and without pass-through columns;
  • oversized output produced by finish(), with and without pass-through columns;
  • draining all fragments from one partition before processing the next partition;
  • result ordering, row count, values, and pass-through alignment;
  • releasing pending partition and result state when the operator is closed early.

The integration test configures a small TsBlock target and uses a table UDF that expands four input rows into 256 variable-width rows. It verifies that every expected row is returned exactly once, payload values are complete, and pass-through columns including null values remain aligned.

Targeted unit test:

mvn test -pl iotdb-core/datanode -am \
  -Dtest=TableFunctionOperatorTest \
  -DfailIfNoTests=false \
  -Dsurefire.failIfNoSpecifiedTests=false

Targeted integration test:

mvn verify -Drat.skip=true -DskipUTs \
  -Dit.test=IoTDBUserDefinedTableFunctionIT#testLargeResultIsSplitWithoutDataLoss \
  -DfailIfNoTests=false \
  -Dfailsafe.failIfNoSpecifiedTests=false \
  -pl integration-test -am \
  -PTableSimpleIT -P with-integration-tests

This PR has:

  • been self-reviewed.
  • added comments explaining non-obvious behavior and limitations.
  • added or updated unit tests.
  • added an integration test.
  • been tested in a test IoTDB cluster.

Key changed/added classes
  • TableFunctionOperator
    • Reuses AbstractOperator to lazily split final table-UDF results after pass-through columns are appended.
    • Keeps queued and retained fragments visible to the operator completion state.
    • Releases pending result and partition references during close.
  • TableFunctionOperatorTest
    • Covers process, finish, pass-through, partition-transition, and early-close paths.
  • LargeResultTableFunction
    • Generates wide integration-test output.
  • IoTDBUserDefinedTableFunctionIT
    • Verifies end-to-end result completeness and pass-through alignment.

@codecov

codecov Bot commented Jul 30, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 88.88889% with 2 lines in your changes missing coverage. Please review.
✅ Project coverage is 43.48%. Comparing base (3a06ed4) to head (ce8dbea).
⚠️ Report is 13 commits behind head on master.

Files with missing lines Patch % Lines
...erator/process/function/TableFunctionOperator.java 88.88% 2 Missing ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master   #18333      +/-   ##
============================================
+ Coverage     43.41%   43.48%   +0.06%     
  Complexity      374      374              
============================================
  Files          5366     5394      +28     
  Lines        382956   385142    +2186     
  Branches      49809    50088     +279     
============================================
+ Hits         166259   167474    +1215     
- Misses       216697   217668     +971     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@sonarqubecloud

Copy link
Copy Markdown

Quality Gate Failed Quality Gate failed

Failed conditions
C Reliability Rating on New Code (required ≥ A)

See analysis details on SonarQube Cloud

Catch issues before they fail your Quality Gate with our IDE extension SonarQube for IDE

Copilot AI 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.

Pull request overview

This PR addresses oversized TsBlock outputs from table-model UDFs by enforcing configured TsBlock size/row limits in TableFunctionOperator, preventing large blocks from propagating through the query pipeline and risking RPC frame overflows.

Changes:

  • Add final-result TsBlock splitting in TableFunctionOperator (after pass-through columns are appended) based on configured TsBlock constraints.
  • Add a unit test covering variable-width outputs with/without pass-through columns.
  • Add an integration-test UDF plus an IT validating large UDF results are returned without loss/duplication and with pass-through alignment under a small TsBlock limit.

Reviewed changes

Copilot reviewed 4 out of 4 changed files in this pull request and generated 3 comments.

File Description
iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/process/tvf/TableFunctionOperatorTest.java Adds unit coverage for splitting variable-width UDF output blocks with optional pass-through columns.
iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/process/function/TableFunctionOperator.java Implements final-result TsBlock splitting and adjusts operator completion/return-size reporting.
integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/udf/IoTDBUserDefinedTableFunctionIT.java Adds end-to-end validation that large table-UDF results are split without data loss and keep pass-through columns aligned.
integration-test/src/main/java/org/apache/iotdb/db/query/udf/example/relational/LargeResultTableFunction.java Introduces a table UDF used by the integration test to generate large variable-width outputs.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines 423 to 426
@Override
public long calculateMaxReturnSize() {
return Math.max(DEFAULT_MAX_TSBLOCK_SIZE_IN_BYTES, properBlockBuilder.getRetainedSizeInBytes());
return maxTsBlockSizeInBytes;
}
Comment on lines +264 to +277
/**
* Splits the final result using the same logical in-memory size accounting as {@link
* TsBlockBuilder}.
*
* <p>Serializing candidate regions to find their exact sizes would write every value into
* temporary buffers, only for the exchange layer to serialize the selected regions again.
* Rebuilding the result with a size-tracking {@link TsBlockBuilder} would avoid that temporary
* serialization, but it would turn the UDF's batched column output into row-by-row,
* column-by-column copies.
*
* <p>Instead, fixed-width values are accounted for directly from their data types, while only the
* retained sizes of variable-width values are inspected. This deliberately estimates the
* in-memory TsBlock size rather than its serialized size because the two representations are not
* equivalent. The resulting regions are views over the original columns and do not copy their
@@ -62,14 +64,15 @@ public class TableFunctionOperator implements ProcessOperator {
private static final long INSTANCE_SIZE =
RamUsageEstimator.shallowSizeOfInstance(AggregationMergeSortOperator.class);

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.

Suggested change
RamUsageEstimator.shallowSizeOfInstance(AggregationMergeSortOperator.class);
RamUsageEstimator.shallowSizeOfInstance(TableFunctionOperator.class);

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.

make TableFunctionOperator extends AbstractOperator, an reuse the functions in that class just like AbstractTableScanOperator.

Copilot AI 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.

Pull request overview

Copilot reviewed 4 out of 4 changed files in this pull request and generated no new comments.

Suppressed comments (2)

iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/process/function/TableFunctionOperator.java:298

  • calculateMaxReturnSize() now returns only maxReturnSize, but calculateMaxPeekMemory() explicitly accounts for Math.max(maxReturnSize, properBlockBuilder.getRetainedSizeInBytes()). This makes the operator’s return-size estimate inconsistent and can underestimate peak output size when the UDF builds wide/variable-width columns (and AbstractOperator may still return a single-row TsBlock larger than maxReturnSize when oneTupleSize > maxReturnSize). Consider restoring the previous max() logic here so memory accounting remains conservative.
  public long calculateMaxPeekMemory() {
    return inputOperator.calculateMaxPeekMemory()
        + Math.max(maxReturnSize, properBlockBuilder.getRetainedSizeInBytes());
  }

  @Override
  public long calculateMaxReturnSize() {
    return maxReturnSize;
  }

iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/process/function/TableFunctionOperator.java:200

  • Splitting via AbstractOperator.checkTsBlockSizeAndGetResult() is based on an average per-row size computed from the whole TsBlock. For variable-width columns with skewed row sizes, this can still produce regions whose retained/serialized size significantly exceeds max_tsblock_size_in_bytes (e.g., one very large row in a block with many small rows lowers the average). This doesn’t match the PR description’s per-position accounting for variable-width columns and may not reliably prevent oversized downstream blocks.
  private TsBlock getNextResultTsBlock() {
    if (retainedTsBlock != null) {
      return getResultFromRetainedTsBlock();
    }
    resultTsBlock = resultTsBlocks.poll();
    return resultTsBlock == null ? null : checkTsBlockSizeAndGetResult();
  }

Copilot AI 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.

Pull request overview

Copilot reviewed 4 out of 4 changed files in this pull request and generated no new comments.

Suppressed comments (1)

iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/process/function/TableFunctionOperator.java:301

  • calculateMaxReturnSize() now returns maxReturnSize only, which can be smaller than properBlockBuilder.getRetainedSizeInBytes() (e.g., many output columns). This can underestimate the operator’s maximum returned block footprint for memory planning, and it regresses the previous behavior that guarded against this by taking Math.max(...).
  @Override
  public long calculateMaxReturnSize() {
    return maxReturnSize;
  }

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.

4 participants