skip to content

A nightly BigQuery load of one 300 GB gzipped CSV from Cloud Storage takes hours. Why, and what would you change?

level: seniorimportance: should knowfreq 45%

answer

  1. parallelism is bounded by readable units
  2. one compression scheme cannot be resumed mid-file
  3. binary block formats keep compression and parallelism
  4. check the job's own duration before blaming the load

basics

~20 s

Gzip is not splittable, so BigQuery must read that file with a single worker — the load is serialised no matter how much capacity exists. Split the data into many files, or use Avro or Parquet, whose internal blocks can be read in parallel.

solid answer

~50 s

A load job scales by reading source files in parallel, and a gzip-compressed CSV cannot be split: the decompressor must run from byte zero, so one worker reads all 300 GB sequentially. CSV also costs the most per byte to ingest — every field is parsed from text and coerced to a type. Two fixes, in order of preference. Produce **Avro or Parquet** instead: they are binary, self-describing, and block-structured, so compressed data still loads in parallel — Google recommends Avro for loading compressed data for exactly this reason. If you are stuck with CSV, emit **many files** (one per partition or shard of the export) rather than one, since parallelism is then bounded by file count; uncompressed CSV is also splittable, so it can beat one gzip file despite the extra bytes. Check `INFORMATION_SCHEMA.JOBS` to confirm the job duration rather than upstream export time is the problem.

code

bash · 5 lines
bash
# Slow: one non-splittable stream, one reader
bq load --source_format=CSV mydataset.events gs://bucket/events/2026-08-20.csv.gz

# Fast: many block-structured files read in parallel
bq load --source_format=PARQUET mydataset.events 'gs://bucket/events/dt=2026-08-20/*.parquet'

go deeper

for a junior

Know that BigQuery loads many files in parallel, and that a single compressed CSV is read by one worker, so file layout drives load time.

for a middle

Explain splittability: why a gzip stream cannot be read from the middle, why block-structured formats like Avro and Parquet can, and why CSV parsing costs more per byte.

for a senior

Diagnose from INFORMATION_SCHEMA.JOBS before changing anything, then prescribe the upstream export layout — many medium files, binary format, dated prefixes — rather than tuning the load job.

for a principal

Own the interface with upstream producers: the file format, file size and partition prefix are a contract, and setting it once removes an entire recurring class of ingestion incidents.

## Where the time goes A BigQuery load job is a parallel read: workers are assigned pieces of the source and decode them concurrently. Anything that prevents the input from being divided collapses that parallelism, and the job degrades to the throughput of one reader. A **gzip stream is not splittable**. Its compression state at any byte depends on every byte before it, so a worker cannot start at the middle of the file. Consequently one gzipped object equals one reader, regardless of size and regardless of available capacity. Three hundred gigabytes read serially is hours. On top of that, **CSV is the most expensive format to ingest**. Every value arrives as text and must be scanned for delimiters and quotes, then parsed and coerced into the column's type. Quoted newlines make it worse, because the reader can no longer treat a newline as a record boundary without tracking quote state. ## The fix, in order **1. Change the format to Avro or Parquet.** Both are binary and self-describing — the schema travels with the data, so there is no type coercion from text and no schema flag to maintain. Both are block-structured: compression is applied per block, so a compressed file remains splittable and workers can read blocks in parallel. Google's guidance is that Avro is the fastest format for loading compressed data for precisely this reason. Parquet, being columnar, is the natural choice when the same files are also read by lake engines or exposed as an external table. **2. If the format cannot change, change the file count.** Parallelism is bounded by the number of independently readable units. One gzip file gives one; a thousand gzip files give a thousand readers. Most export tools can shard output — do that, and keep the shards in the hundreds of megabytes rather than kilobytes, since thousands of tiny files trade one bottleneck for per-file overhead. **3. Consider dropping the compression.** Uncompressed CSV is splittable, so a single large uncompressed file loads in parallel while a single gzipped one does not. You pay more Cloud Storage egress-free bytes to read and more storage on the staging bucket, but the wall-clock difference on a serialised load is usually decisive. Sharded-and-compressed still beats both. ## Confirm before you act Do not assume the load job is the slow part. `INFORMATION_SCHEMA.JOBS` gives creation, start and end times per load job, along with bytes and errors: ```sql SELECT job_id, creation_time, start_time, end_time, TIMESTAMP_DIFF(end_time, start_time, MINUTE) AS minutes, total_bytes_processed, error_result FROM `region-us`.INFORMATION_SCHEMA.JOBS WHERE job_type = 'LOAD' AND creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 7 DAY) ORDER BY minutes DESC; ``` If the job itself runs for minutes but the pipeline takes hours, the upstream export or the file-arrival wait is the real target. If the job runs for hours on one file, you have confirmed the serialisation. ## Secondary causes worth ruling out - **Slot contention.** Batch loads run on a free shared pool by default. Under heavy contention they can queue. Routing loads to a reservation gives them dedicated capacity — worth doing when a load has an SLA. - **Location mismatch.** The source bucket must be colocated with the dataset; a badly placed bucket is a hard failure rather than slowness, but it is worth checking when a pipeline is newly built. - **Autodetect on a huge CSV.** Schema inference samples the input; pinning an explicit schema removes both the sampling cost and the risk of a type flipping between runs. - **A single wide row.** Very large individual records can throttle a reader independent of the file layout. ## The durable shape The pattern that avoids this class of problem entirely: the upstream job writes **many medium-sized Parquet or Avro files, partitioned by date**, into a dated prefix; the load job targets `gs://bucket/table/dt=YYYY-MM-DD/*.parquet` with `WRITE_TRUNCATE` against that partition. Parallel by construction, cheap to parse, atomic per partition, and idempotent on rerun.

  • Why does an Avro file compressed with deflate still load in parallel when a gzipped CSV does not?
    Avro compresses per data block and records block boundaries in the file, so a reader can seek to a block and decompress it independently. A gzipped CSV is a single compression stream with no block index, so decoding any byte requires decoding everything before it — one file, one reader.
  • Would splitting the CSV into 100,000 tiny gzip files fix it?
    It removes the serialisation but introduces per-file overhead: each object costs a listing entry, a connection and a task. Aim for a few hundred megabytes per file — enough files to saturate parallelism, large enough that setup cost is amortised. Thousands of very small files is a known way to make a load slower again.
  • How would you tell whether the load job or the upstream export is the slow step?
    Query INFORMATION_SCHEMA.JOBS for the load job's start and end times. If the job itself runs for hours, the source layout is the problem. If it runs for minutes but the pipeline is late, the export, the file-arrival wait or the scheduler is the target, and changing the file format will not help.

saying these in an interview costs you the question

  • Assuming BigQuery parallelises inside a single gzip stream
  • Blaming BigQuery capacity for a load bottlenecked by one file
  • Thinking compressed always loads faster than uncompressed
  • Splitting into hundreds of thousands of tiny files as the fix
  • Never checking the load job's actual duration before changing the pipeline

context