Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Batch size and parallelism for loads

Pipelines & scenarios · Core Concepts Not Yet Covered

Batch size and parallelism for loads

Mediumpipelines-65
batch-sizethroughputbulk-loadrate-limitssinks

Question

How do you decide batch size when loading data into a database or API sink?

Solution

The best batch size is the largest one that stays within the sink's limits and keeps retries cheap. Too small wastes time on overhead, and too large makes failures painful.

The two ends

  • Too small: sending one row per request or per insert has a fixed cost each time (network round trip, connection handling, transaction commit). Inserting 1 million rows one at a time into a database can take hours. The same rows in batches of 5,000 take minutes.
  • Too large: one huge request can time out, use too much memory on your side or the sink's, hit payload or transaction size limits, and, if it fails at 95 percent, you redo everything. Long transactions on a database also hold locks and bloat logs.

Know the sink's limits

Every destination has constraints: an API may accept at most 1,000 records or a few MB per call, with a rate limit per second. A database has maximum transaction sizes and connection counts. A message bus has a maximum message size. Read the documentation, and set your batch size below the limits with some margin.

Use the bulk path

Row-by-row inserts are almost always the slowest option. Warehouses and databases have bulk loading paths: COPY in Postgres and Snowflake, load jobs in BigQuery, COPY INTO in Databricks, or bulk APIs for SaaS tools. They load files or large batches far faster, and often more cheaply.

Parallelism

Several workers sending batches at once can raise throughput, up to what the sink can take. Beyond that, you only get slowdowns, errors and angry database administrators. Increase parallelism step by step while watching error rates, latency and the sink's load. For APIs, respect the rate limit across all workers together, not per worker.

Tune by measuring

Try batch sizes such as 500, 1,000, 5,000 and 20,000, measure rows per second and failure behaviour, and choose a value in the flat part of the curve. Do not pick the single fastest one if a smaller size is nearly as fast and retries cost less.

Make retries safe

With batches, a retry resends the whole batch, so writes should be idempotent (upsert by key). Keep batches independent, so one failure does not stop the others, and record which batches are done.

PreviousNext