Skip to content

[Bug]: Window operator fails when agg Flush returns split vectors (>8192 rows): "does not support sending split result of window function" #25813

Description

@heni02

Is there an existing issue for the same bug?

  • I have checked the existing issues.

Branch Name

main

Commit ID

f8aaa25

Other Environment Information

- Hardware parameters: N/A(逻辑 bug,与机器无关)
- OS type: Linux / Darwin
- Others: 单机或 multi-CN 均可复现;触发条件是 Window 输入行数 > AggBatchSize(8192)

Actual Behavior

执行下面这条带 ROW_NUMBER + SUM OVER 的查询时报错:

Error: -16 ... [ERROR] [execute_sql.py:51, _executesql_impl]
(20101, 'internal error: the Window operator currently does not support sending split result of window function.')

Failing SQL(线上/回归实际触发)

WITH w AS (
WITH order_dim AS (
  SELECT vbeln, MAX(zsal_rep) AS zsal_rep, MAX(vkgrp) AS vkgrp, MAX(vkbur) AS vkbur, MAX(zsal_det) AS zsal_det
  FROM dwd_dcp.dwd_s4_ztbw_qcbb GROUP BY vbeln
),
filtered_acdoca AS (
  SELECT a.rbukrs, a.belnr, a.gjahr, a.zuonr, a.kunnr, a.wsl, a.prctr, a.vkbur_pa
  FROM dwd_dcp.dwd_s4_acdoca a
  WHERE a.drcrk = 'H' AND (a.racct LIKE '1122%' OR a.racct LIKE '2205%')
    AND a.blart IN ('DZ','DW','DM') AND (a.xtruerev != 'X' OR a.xtruerev IS NULL)
    AND a.gjahr IN ('2022', '2023', '2024', '2025')
    AND a.rbukrs IN ('1000', '1100', '1200', '1300', '2000', '2100', '2200', '3000')
)
SELECT
  a.rbukrs, a.belnr, a.gjahr, a.zuonr, a.kunnr, a.wsl,
  COALESCE(c.vkgrp, od.vkgrp) AS vkgrp,
  COALESCE(tgrp.bezei, od.zsal_rep) AS sales_group_desc,
  e.cbukrs, e.cbuktx,
  COALESCE(f1.cprctr, f2.cprctr) AS cprctr,
  g.cpatnr_txtlg, NULLIF(h.comp_head, '/') AS comp_head,
  m_code.dept_id, m_desc.dept_name,
  ROW_NUMBER() OVER (PARTITION BY a.rbukrs, a.kunnr ORDER BY a.wsl DESC) AS rn,
  SUM(a.wsl) OVER (PARTITION BY a.rbukrs) AS bukrs_wsl_total
FROM filtered_acdoca a
LEFT JOIN order_dim od ON a.zuonr = od.vbeln
LEFT JOIN dwd_dcp.dwd_s4_vbak c ON a.zuonr = c.vbeln
LEFT JOIN dwd_dcp.dwd_s4_tvgrt tgrp ON c.vkgrp = tgrp.vkgrp
LEFT JOIN dwd_dcp.dwd_s4_tvkbt tbur ON c.vkbur = tbur.vkbur
LEFT JOIN dwd_dcp.dwd_bw_ztbpc002_com e ON a.rbukrs = e.sbukrs
LEFT JOIN dwd_dcp.dwd_bw_ztbpc002_prc f1
  ON f1.sysid = 'CN1' AND f1.bukrs = a.rbukrs AND f1.sprctr = a.prctr AND f1.bukrs != ''
LEFT JOIN dwd_dcp.dwd_bw_ztbpc002_prc f2
  ON f2.sysid = 'CN1' AND (f2.bukrs = '' OR f2.bukrs IS NULL) AND f2.sprctr = a.prctr
LEFT JOIN dws_dcp.dws_bpqx g ON g.sysid = 'CN1' AND g.customer = a.kunnr
LEFT JOIN dwd_dcp.dwd_s4_bp001 h ON a.kunnr = h.partner
LEFT JOIN staging_db.sales_office_mapping m_code
  ON c.vkbur = m_code.sales_office_code
LEFT JOIN staging_db.sales_office_mapping m_desc
  ON tbur.bezei = m_desc.sales_office_desc
)
SELECT * FROM w
WHERE 1 = 1
LIMIT 500;

触发点主要是 JOIN 后的两个窗口函数(尤其 SUM(a.wsl) OVER (...)),当进入 Window 的行数 > AggBatchSize(8192) 时 Flush() 返回多段 vector,Window 直接报错。

Expected Behavior

Window 算子应能正确处理 Flush() 返回的多段 result vector(与 Group 算子已支持的 split result 行为一致),查询成功返回结果。

Steps to Reproduce

Minimal repro(不依赖业务表)

DROP TABLE IF EXISTS t_win_split;
CREATE TABLE t_win_split (g INT, v DECIMAL(20,2));

-- 插入 >8192 行,迫使 SUM OVER 的 agg state 分块
INSERT INTO t_win_split
SELECT
  (result % 8) AS g,
  CAST(result AS DECIMAL(20,2)) AS v
FROM generate_series(1, 10000) g(result);

-- 触发:Window 对整表/大 batch 做 SUM OVER,Flush 返回 >1 个 vector
SELECT
  g, v,
  SUM(v) OVER (PARTITION BY g) AS part_sum,
  ROW_NUMBER() OVER (PARTITION BY g ORDER BY v DESC) AS rn
FROM t_win_split
LIMIT 100;

若环境没有 generate_series,可用任意方式插入 ≥9000 行后再跑同一条 SELECT

(上面 Failing SQL 即业务侧完整复现;若无金盘表,用本节 Minimal repro 即可。)

Additional information

Root cause(代码位置)

  1. pkg/sql/colexec/window/window.goFlush() 后只接受单个 vector:
vecs, err := ctr.batAggs[idx].Flush()
...
if len(vecs) > 1 {
    return moerr.NewInternalErrorNoCtx(
        "the Window operator currently does not support sending split result of window function.")
}
  1. 聚合执行器按 AggBatchSize = 8192 分块(pkg/sql/colexec/aggexec/aggState.go)。
    SUM / decimal SUM fast path(sum_decimal_fast.go / sumavg2.go)的 Flush() 返回 len(exec.state) 个 vector。当 group/行数 >8192 时 len(vecs) > 1

  2. 对比:Group 已能消费多段 Flush 结果(把各段挂到对应 output batch),Window 尚未对齐。

Suggested fix

  • Window:对 Flush() 多段结果做 UnionBatch 合并,或按 chunk 拆成多个 output batch(与 Group 对齐)。
  • 补充回归:SUM/AVG/... OVER + 输入行数 >8192(含 PARTITION BY / 无 PARTITION BYreceiveAll 路径)。

Workaround

将窗口改写为 JOIN + GROUP BY(例如 SUM OVER (PARTITION BY g) → 先 GROUP BY g 再 join 回明细),可绕过 Window 算子。

Metadata

Metadata

Assignees

Labels

kind/bugSomething isn't workingseverity/s0Active / top priority for current sprint. Owner has committed to working on it now.

Type

Projects

No projects

Milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions