Skip to content

ingest: split oversized batches by halving, and don't split TOO_MANY_PARTS caused by merge backpressure #657

Description

@EricAndrechek

Area: ingest — batch insert retries (follow-up to #619)

#619 made the ingest worker split a multi-row batch that ClickHouse refuses with TOO_MANY_PARTS (252) or MEMORY_LIMIT_EXCEEDED (241) (chconn.Splittable), using the same row-by-row isolation as a rejected batch. The motivation is sound: a batch spanning more partitions than max_partitions_per_insert_block fails with 252, is re-formed the same way on redelivery, and never lands. Two problems remain:

  1. TOO_MANY_PARTS has two causes, and only one is helped by splitting.

    • Too many partitions in one INSERT ("Too many partitions for single INSERT block (more than N)"): caused by the batch, so a smaller insert fixes it.
    • Too many active parts because merges are behind ("Too many parts (N) … Merges are processing significantly slower than inserts"): caused by the table, and every extra insert makes it worse.

    The code treats both as splittable. Under merge backpressure, if the first single-row insert gets through (a merge just finished), isolation carries on and inserts the rest one row per request. Each of those creates a new part, which feeds the condition ClickHouse is reporting. Only the partitions-per-insert case should split; the merge-backpressure case should stay a plain backoff.

  2. Row-by-row is the wrong granularity for a size-driven split. A 10,000-row batch across 150 day-partitions becomes 10,000 inserts and 10,000 tiny parts, which can itself trip the parts limit. The same applies to MEMORY_LIMIT_EXCEEDED, which can also be server-wide memory pressure from other queries rather than this batch's size.

Expected:

  • Tell the two TOO_MANY_PARTS causes apart by the exception message. Split only the partitions-per-insert case, and back off on merge backpressure.
  • For the splittable cases, halve recursively: retry each half, recurse only into halves that fail again. That lands in about log₂(n) requests and keeps parts large. If a single row still fails the same way, the server is the cause: back off, as today.
  • Keep row-by-row isolation for Rejected batches only.
  • Tests:
    • merge backpressure: no split, backoff;
    • a batch spanning 150 partitions: halves and lands in a few requests, with the request count asserted;
    • a server-wide memory limit: backs off after about log₂(n) attempts;
    • a single-row batch: unchanged.

Related: #619, #649.

Activity

  1. EricAndrechek commented on Sep 26, 2026

    @EricAndrechek
    MemberAuthor

    A related item, for later: a change tracked upstream in chtypes would let it answer the batch-shaped 252 ("too many partitions for a single INSERT block") before the INSERT. It would evaluate the table's PARTITION BY per row with ClickHouse's own code, refuse a body over max_partitions_per_insert_block as the server does, and return per-row partition ids so a producer can split by partition up front. Once that exists, a 252 that reaches the worker means merge backpressure only. It needs chtypes' next ABI revision, so it is not imminent; this issue's message-based split stands until then.

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions