skip to content

Why can one 20 GB gzipped CSV file in S3 make a Redshift Spectrum query slow?

level: middleimportance: nice to knowfreq 35%

answer

  1. parallelism needs somewhere to cut the file
  2. not every compression codec can be entered midway
  3. one object can become one worker
  4. the opposite mistake is a million tiny objects
  5. look at splits and request parallelism

basics

~20 s

Gzip is not splittable, so the whole file must be decompressed sequentially by a single Spectrum reader. No matter how large the cluster, one worker does all the work, and because CSV is row-based every column is decoded too.

solid answer

~50 s

Spectrum parallelises by handing readers independent pieces of the data. Parquet and ORC split naturally at row-group and stripe boundaries, and uncompressed or bzip2 text can be split at record boundaries. A **gzip** stream cannot be entered partway through, so one file means one reader decompressing 20 GB serially while the rest of the fleet idles. CSV compounds it: being row-based, that single reader also decodes all forty columns to serve a two-column query. The opposite failure is just as real — a million 200 KB files spend all their time on per-object requests and listing rather than scanning. Aim for a modest number of reasonably large, evenly sized, splittable columnar files per partition. Check `files`, `splits`, `avg_request_parallelism` and `max_request_duration` in `SVL_S3QUERY_SUMMARY`: one long-running request beside many fast ones is the signature of an unsplittable or oversized file.

code

sql · 7 lines
sql
SELECT files,
       splits,
       avg_request_parallelism,
       max_request_duration,
       s3_scanned_bytes
FROM   svl_s3query_summary
WHERE  query = pg_last_query_id();

go deeper

for a junior

Know that Spectrum reads S3 files in parallel, and that a single large gzipped file cannot be divided, so only one reader handles it however big the cluster is.

for a middle

Explain splittability by format — row groups in Parquet, record boundaries in plain text, a continuous stream in gzip — and why both very large and very tiny files hurt for different reasons.

for a senior

Show you would measure it: files, splits, avg_request_parallelism and max_request_duration, then fix the writing pipeline with compaction and even file sizing rather than tuning the query.

for a principal

Own the file-layout contract with upstream producers — target sizes, format, compaction cadence — because query teams cannot fix a layout they do not control.

## How Spectrum gets its parallelism The Spectrum scan fleet is elastic, but it can only apply that elasticity if the data can be divided. Work is handed out as **splits** — independent ranges a reader can process without seeing the rest. Where those boundaries come from depends entirely on the file format and its compression. - **Parquet and ORC** are internally chunked into row groups and stripes, each self-describing with its own footer metadata. Any reader can jump to any row group. These formats split well by construction, even when their internal pages are compressed, because the compression is applied per column chunk rather than across the whole file. - **Plain text (CSV, JSON lines), uncompressed** can be split at record boundaries: a reader starts at a byte offset, skips to the next newline, and proceeds. - **bzip2-compressed text** is block-structured and remains splittable. - **gzip-compressed text** is a single continuous DEFLATE stream. There is no way to start decoding at byte 10,000,000 without having decoded everything before it. **One gzip file is one split, and therefore one reader.** ## What that does to a 20 GB gzipped CSV The query's entire S3 phase collapses onto a single worker. It must: 1. stream and decompress 20 GB serially, 2. parse every record, 3. decode **every column**, because CSV has no column-level structure, and 4. only then apply your filter and projection. Adding nodes to the Redshift cluster does not help — the bottleneck is a single-threaded decompression of one object. Nor does the filter help the scan cost: the bytes were read regardless. ## The mirror-image failure: too many small files The intuitive fix — "split it into lots of files" — overshoots easily. Every object carries fixed costs: a listing entry, a request round-trip, format metadata to parse, a reader assignment. When a partition contains hundreds of thousands of tiny objects, those fixed costs dominate and throughput collapses again, this time from overhead rather than serialisation. Streaming ingestion that writes one file per micro-batch produces exactly this shape, which is why lake pipelines almost always run a compaction step. ## Skew Even with a splittable format, badly *uneven* files hurt. If one partition holds a single 30 GB file and forty 50 MB files, the readers assigned to the small files finish immediately and the query waits on the long pole. The goal is not just "enough splits" but **evenly sized** splits, so all readers finish at roughly the same time — the same straggler logic that governs any parallel scan. ## What to aim for There is no single magic number, and it depends on partition size and query shape, but the shape of the target is consistent: - **Columnar and splittable**: Parquet or ORC, with internal compression such as Snappy or zstd. - **Substantially larger than trivial**: files in the tens-to-hundreds-of-megabytes range rather than kilobytes, and not so large that a single file becomes the straggler. - **Enough files per partition to occupy the readers** — file count comfortably above the number of parallel readers the query can use, which relates to the cluster's slice count. - **Evenly sized**, produced by a compaction job rather than by whatever the ingestion batch happened to emit. ## Diagnosing it `SVL_S3QUERY_SUMMARY` is the instrument. Its `files` and `splits` columns say how divisible the data actually was — `splits` barely exceeding `files` for a supposedly columnar table is a red flag. `avg_request_parallelism` shows how many readers were actually busy. `max_request_duration` beside a much smaller average is the straggler signature: one file, one reader, everyone else waiting. ```sql SELECT files, splits, avg_request_parallelism, max_request_duration, s3_scanned_bytes FROM svl_s3query_summary WHERE query = pg_last_query_id(); ``` If the numbers show one split for a large object, the fix is upstream in the pipeline that wrote it, not in the SQL.

  • Which S3 file layouts split well for Spectrum, and which do not?
    Parquet and ORC split at row-group and stripe boundaries and are the right default. Uncompressed text splits at record boundaries, and bzip2 text is block-structured so it splits too. Gzip-compressed text is a single continuous stream and cannot be split, so each such file is processed by exactly one reader.
  • Why is a million small Parquet files also a problem?
    Every object costs a listing entry, a request round-trip, footer metadata parsing and a reader assignment. Below a certain size those fixed costs exceed the useful scan work, so throughput collapses from overhead. Streaming pipelines that write one file per micro-batch hit this routinely and need a compaction job.
  • How do you spot a straggler file from Redshift's own system views?
    Query `SVL_S3QUERY_SUMMARY` for the statement and compare `max_request_duration` against the average, alongside `avg_request_parallelism` and the `files`/`splits` counts. A single long request while parallelism is low means one oversized or unsplittable object is holding the query open while the other readers sit idle.

saying these in an interview costs you the question

  • Assuming gzip and Parquet parallelise the same way
  • Thinking a bigger cluster fixes a single unsplittable file
  • Splitting data into millions of tiny objects to increase parallelism
  • Ignoring uneven file sizes because the format is columnar
  • Blaming S3 throughput when the cause is one long-running reader

context