skip to content

What delivery and schema semantics should you expect from the Snowflake and BigQuery Kafka Connect sinks, and how do they differ from a file-based S3 sink?

level: seniorimportance: nice to knowfreq 30%

answer

  1. S3 = files to catalog; warehouse = rows to tables
  2. Snowflake: Snowpipe Streaming, VARIANT or schematized columns, offset tracking = EOS
  3. BigQuery: Storage Write API = exactly-once, auto-create/evolve schema
  4. warehouse rows instantly queryable
  5. watch streaming quotas/cost

basics

~20 s

Warehouse sinks (Snowflake, BigQuery) write rows directly into managed tables rather than files in a bucket. They use the warehouse's streaming-ingest APIs, can offer exactly-once or near-exactly-once via offset tracking, and auto-map record schemas to table columns — whereas the S3 sink just writes batched files you must catalog separately.

solid answer

~50 s

A **file-based S3 sink** writes batched objects (Parquet/Avro/JSON) into a bucket; you then register them in a catalog/external table to query. **Warehouse sinks** load rows into managed tables. The **Snowflake sink** historically used Snowpipe (micro-batch, near-real-time) and now **Snowpipe Streaming** for low-latency row inserts; it lands data in VARIANT columns (`RECORD_CONTENT`/`RECORD_METADATA`) or, with schematization enabled, maps fields to typed columns, and tracks Kafka offsets to provide **exactly-once** ingestion. The **BigQuery sink** (e.g., Confluent/WePay) writes via the legacy streaming insert or the **Storage Write API** (which supports exactly-once with stream commits), auto-creates/updates tables from the record schema, and supports upsert/delete for CDC on newer versions. Key differences from S3: rows are immediately queryable (no separate cataloging), schema is mapped to warehouse columns (auto-evolution possible), and exactly-once relies on the warehouse API's offset/commit semantics rather than deterministic file naming. Trade-offs: warehouse streaming-ingest cost, quota limits, and schema-evolution compatibility.

go deeper

for a junior

Know warehouse sinks load rows straight into Snowflake/BigQuery tables, unlike S3 which writes files.

for a middle

Explain Snowpipe Streaming vs file staging and that warehouse sinks map schema to columns.

for a senior

Compare exactly-once mechanisms (offset tracking / Storage Write API) vs S3 file naming, and CDC upsert support.

for a principal

Decide between open-lake (S3/Iceberg) and direct-warehouse sinks per cost, lock-in, latency, and multi-engine needs.

## Two shapes of sink 1. **File/object sinks** (S3, GCS, ADLS): write **batched files** into a bucket. The data isn't a queryable table until you point an **external table / catalog** (Glue, Hive metastore, Iceberg) at it. Delivery semantics come from **deterministic offset-based file naming** (see the S3 sink). You own file layout, compaction, and cataloging. 2. **Warehouse sinks** (Snowflake, BigQuery): write **rows directly into managed tables** via the warehouse's ingest API. The warehouse owns storage and indexing; data is **immediately queryable**. ## Snowflake sink - **Ingest path**: originally **Snowpipe** (micro-batch file staging behind the scenes, near-real-time); modern versions use **Snowpipe Streaming** for **low-latency, row-level** inserts without intermediate files. - **Schema**: by default lands two columns — `RECORD_CONTENT` (the record, often as **VARIANT** semi-structured JSON) and `RECORD_METADATA` (topic, partition, offset, timestamp, key). With **schematization** enabled, it parses the record schema and maps fields to **typed columns**, and can **auto-evolve** the table when new fields appear. - **Delivery**: it **tracks Kafka offsets** per channel/partition in Snowflake, enabling **exactly-once** ingestion — on replay it recognizes already-committed offsets and skips duplicates. ## BigQuery sink - **Ingest path**: legacy **streaming inserts** (tabledata.insertAll) or the newer **Storage Write API**, which supports **stream-level commit** and thus **exactly-once** semantics. - **Schema**: can **auto-create** the destination table and **update schema** from the record/Schema Registry definition; nested Avro/Protobuf maps to BigQuery RECORD/STRUCT. - **CDC**: recent connector versions support **upsert and delete** (merge by key) so Debezium streams materialize current state directly in BigQuery. - **Quotas/cost**: streaming ingestion has per-table quotas and per-GB streaming cost to plan for. ## Key contrasts with S3 - **Queryability**: warehouse rows are instantly queryable; S3 files need cataloging/external tables. - **Schema mapping**: warehouse sinks map record fields to **typed columns** and can auto-evolve; S3 stores whatever file format you chose, schema enforced at read time. - **Exactly-once mechanism**: warehouse sinks lean on **API offset/commit semantics** (Snowflake offset tracking, BigQuery Storage Write API stream commits); the S3 sink relies on **deterministic file naming + atomic PUT**. - **Cost/limits**: warehouse streaming-ingest costs and quotas vs. S3's cheap object storage but extra catalog/compaction overhead. ## Practical guidance Use warehouse sinks when consumers live in that warehouse and want instant SQL with minimal ops; use the S3/Iceberg path when you want an open, engine-agnostic, cheap lake with replay and multi-engine access. Many architectures do both: land bronze in S3/Iceberg and also stream hot tables into the warehouse.

  • How does the Snowflake sink achieve exactly-once where the S3 sink uses file naming?
    The Snowflake sink tracks the committed Kafka offset per channel/partition inside Snowflake. On replay it compares incoming offsets to what's already committed and skips duplicates, giving exactly-once at the row level — rather than the S3 sink's deterministic offset-encoded filenames plus atomic object overwrite.
  • What lets the BigQuery sink provide exactly-once that the older streaming insert path didn't?
    The BigQuery Storage Write API supports application-created streams with explicit commit/offset semantics, so a batch is committed exactly once even on retries. The legacy tabledata.insertAll path was at-least-once and relied on insertId best-effort dedup within a short window.

saying these in an interview costs you the question

  • Saying warehouse sinks write files to a bucket like S3 — they insert rows via warehouse ingest APIs.
  • Assuming all warehouse sinks are exactly-once regardless of ingest path (legacy BigQuery streaming inserts were at-least-once).
  • Forgetting streaming-ingest quotas/costs.
  • Thinking S3 files are immediately queryable without a catalog/external table.

context