When should analytical data stay in external tables on object storage instead of being loaded into the engine?
answer
- who wrote the files decides what can be skipped
- how often is this data actually queried
- copies are the price of speed
- hot data loaded, cold history in place
basics
~20 sKeep data external when it is queried rarely, is large and cheap to retain, is still raw or evolving, or must stay readable by other engines. Load it when queries are frequent, latency-sensitive or highly concurrent, because loading buys layout control, statistics, clustering and caching.
solid answer
~50 sAn external table lets an analytical engine read files in object storage in place: no copy, no load job, and other engines keep their access to the same bytes. That is the right shape for cold or rarely-queried history, for a raw landing zone you explore occasionally, for one-off backfills, and for data another team owns and publishes. It stops being right once queries get frequent or interactive. Reading in place means the engine did not choose the file sizes, the sort order, the encodings or the statistics, so pruning is limited to whatever the directory layout and file metadata allow, and small-file sprawl turns into per-file overhead on every scan. Loading into engine-managed storage buys layout control, clustering, richer statistics and caching, which usually turns into a large reduction in bytes scanned. The common design keeps both: hot marts loaded, cold history external, and a union view over the boundary.
code
sql · 12 lines-- read in place: no copy, other engines keep access
CREATE EXTERNAL TABLE events_ext (
event_id BIGINT,
user_id BIGINT,
event_ts TIMESTAMP
)
STORED AS PARQUET
LOCATION 's3://raw-zone/events/';
-- load a copy the engine owns, lays out and keeps statistics on
CREATE TABLE events AS
SELECT * FROM events_ext WHERE event_ts >= DATE '2026-07-01';go deeper
Know that an external table is a definition over files the engine reads in place, with no copy, and that loading makes a second copy the engine controls.
Explain what the engine loses by not owning the files: coarse pruning, thin statistics, per-file overhead, no clustering or caching. Be able to say what a load actually buys.
Make the call with numbers — query frequency, bytes scanned per query, file count and format — and describe the hot-loaded plus cold-external hybrid, including how you keep the boundary consistent and prunable.
Own the placement policy across the platform: what stays open for other engines, what gets copied, who pays for each copy, and how you stop copies proliferating into datasets that disagree.
## What an external table actually is An external table is a table definition — a name, a column list, a location, a file format — where the engine stores **only the definition** and reads the data files from object storage on demand. Nothing is copied. The files remain in whatever format they were written in, owned by whoever wrote them, and readable by any other engine pointed at the same prefix. Most analytical engines support some version of this, along with a partitioning convention where directory names encode a column value so the engine can skip whole directories. The contrast is a **managed** (internal) table, where you run a load that rewrites the data into the engine's own layout. The engine then controls file sizes, sort order, encodings, statistics and caching. ## What you give up by reading in place The engine did not write the files, so most of its performance levers are unavailable: - **Pruning is only as good as the layout.** If the files are partitioned by date and you filter on date, you skip directories. Filter on anything else and you scan everything, because there is nothing else to prune with. A managed table can be clustered or sorted so that many predicates prune. - **Statistics are thin.** The engine can use whatever the file format's own metadata offers plus directory structure. It usually lacks the table-level statistics it maintains for its own tables, so join order and memory estimates get worse — which shows up as a bad plan on a big join, not just a slower scan. - **Small files hurt disproportionately.** Every file costs a listing entry, a request, and a fixed amount of open-and-read overhead. A partition made of ten thousand tiny files from a streaming writer can cost more in overhead than in bytes. - **No result reuse from the engine's own caching layers**, and no clustering maintenance, because the engine does not own the data. - **The definition can drift from reality.** Files change under you; the declared column list is a claim, not a constraint, so a producer's change surfaces as wrong results rather than a rejected load. ## What you get in exchange - **No copy and no load pipeline.** The data is queryable as soon as it lands, and you are not paying to store a second copy or to keep it in sync. - **Other engines keep reading it.** This is the decisive argument when the same dataset feeds a batch engine, a notebook and a warehouse. - **Cheap retention.** Object storage is the cheapest tier you have. Years of history you touch twice a quarter should not be sitting in engine-managed storage. - **Raw fidelity.** Nothing was reshaped on the way in, which matters for reprocessing and audit. ## The decision, in practice Ask how often the data is queried and how much each query reads. A dataset scanned by a handful of ad-hoc queries a month is external. A dataset behind a dashboard fifty people refresh all morning is loaded — the accumulated scan cost of reading unclustered files repeatedly dwarfs the one-time cost of loading, and the latency is the difference between an interactive product and a slow one. Then ask who else needs it. If another engine must read the same table, external (or an open lakehouse table) is the way to avoid a copy that will eventually disagree with the original. Then ask whether the shape is stable. Raw, evolving, semi-structured data belongs external until you know what you want from it; once the questions are stable, model and load the answer. ## The hybrid that most platforms land on Hot recent data lives in managed tables, older partitions are aged out to external files, and a view unions the two so consumers see one table. This keeps the working set fast and the long tail cheap. The pieces to get right are the boundary predicate (queries must be able to eliminate one side entirely), consistency during the move, and making sure the view does not defeat pruning on either side. ## Diagnosis worth naming When an external-table query is slow or expensive, work through: how many files did it open, how many bytes did it read versus the table's total size, and did any partition elimination happen? The three usual causes are no usable partition predicate, thousands of small files, and a text or row-oriented format where a columnar one would have let the engine read only the columns the query asked for. Fixing the file layout — compacting to larger files and rewriting as columnar with a sensible partition scheme — usually beats moving to a bigger cluster.
- An external-table query reads far more data than the filter suggests it should. What do you check first?Whether any partition elimination happened. If the files are laid out by date and the query filters on something else — or wraps the partition column in a function — the engine has nothing to skip with and reads everything. Next check file count and format: thousands of small files or a row-oriented text format both force full reads.
- How do you keep a hot-loaded plus cold-external hybrid table usable for consumers?Expose a single view that unions the two, split on a boundary predicate the optimizer can use to eliminate one side entirely. Automate the aging job, make the cutover to external idempotent so a partition is never missing or duplicated during the move, and monitor that queries still prune on both sides.
- Does storing the external files in a columnar format instead of JSON change the decision?It moves the line substantially. Columnar files let the engine read only the referenced columns and use per-chunk statistics to skip data, so external reads get much closer to managed-table cost. Row-oriented text formats force a full parse of every record and are the strongest argument for loading.
saying these in an interview costs you the question
- Assuming external tables are always cheaper because storage is cheaper
- Expecting the engine to cluster or index files it does not own
- Ignoring per-file overhead when thousands of small files exist
- Treating the declared column list as enforced on the files
- Loading everything by default without measuring query frequency