"Streaming inserts" into BigQuery used to mean one thing: calling tabledata.insertAll with a batch of JSON rows and watching them become queryable within seconds. Today the name covers two different APIs. Google's documentation now calls the old method the BigQuery Storage Write API (REST), noting that it was "previously known as the legacy tabledata.insertAll method", and recommends the Storage Write API (gRPC) for new projects. They land rows in the same tables but promise different things about duplicates, atomicity and failure handling, and they are billed differently.
This article explains both paths from first principles: what happens to a row between your process and a query result, how deduplication works on each path and where it stops working, how committed streams and offsets give exactly-once writes, when pending streams replace load jobs, and how change-data-capture upserts work. A worked design, failure modes and a checklist follow. For BigQuery fundamentals see Google BigQuery; for table layout see BigQuery Partitioning and Clustering.
Two APIs behind one name
Both APIs exist because load jobs have minutes of latency and per-table job quotas. Streaming trades that for row-level freshness: once a write is acknowledged, a query can see the row. The two APIs differ in what an acknowledgement means and what a retry does.
| Property | Storage Write API (REST), formerly insertAll | Storage Write API (gRPC) |
|---|---|---|
| Transport and format | HTTPS, JSON rows | gRPC bidirectional stream, protobuf rows (Arrow is accepted in some modes) |
| Duplicate handling | Optional insertId; best-effort, roughly one-minute window | At-least-once on the default stream; exactly-once on committed streams with offsets |
| Atomic multi-row commit | No | Yes, with pending streams |
| Partial failure | Per-row errors inside a successful HTTP response | Per-request; a bad row fails its append |
| CDC upserts | No | Yes, default stream with _CHANGE_TYPE |
| Google's recommendation | Not recommended for new projects | Recommended for new projects |
Pricing is also different: the gRPC API is billed by ingested bytes with a monthly free allowance, while the REST method has its own per-volume rate. Take current rates from the BigQuery pricing page; Google cites lower pricing as one reason it recommends the gRPC API.
The REST method, formerly insertAll
The REST method is simple, which is why it is still everywhere. You post up to a few hundred rows per request, each optionally carrying an insertId. The response is HTTP 200 even when some rows failed, so the first rule of using it is to read insertErrors on every response. Each entry gives the index of the failed row and a reason. Two flags, skipInvalidRows and ignoreUnknownValues, let valid rows through and drop unknown fields; both are convenient and both hide data problems.
from google.cloud import bigquery
client = bigquery.Client()
TABLE = "proj.analytics.events"
def write_batch(events):
rows = [e.to_row() for e in events]
ids = [e.event_id for e in events] # stable ID from the producer, not uuid4() at send time
errors = client.insert_rows_json(TABLE, rows, row_ids=ids)
if errors: # list of {"index": i, "errors": [...]}
bad = {err["index"] for err in errors}
dead_letter([events[i] for i in bad], errors)
retry_later([e for i, e in enumerate(events) if i not in bad and needs_retry(errors, i)])The insertId is the only duplicate protection on this path. BigQuery uses it to try to drop a repeated row that arrives within about one minute of the first; the documentation is explicit that this is best-effort and duplicates may still appear. Two consequences follow. The ID must be derived from the event itself, so that a retry carries the same ID; generating a fresh UUID per attempt disables deduplication entirely. And any retry that happens more than a minute later, for example after a process restart replays a queue, will produce a duplicate. If correctness matters, plan for read-side deduplication anyway.
Two partitioning details catch people. In an ingestion-time partitioned table, freshly streamed rows show NULL in _PARTITIONTIME until BigQuery assigns them, typically within minutes and in rare cases up to 90 minutes, so a query filtered on today's partition can miss rows that a full scan sees. And when you target a specific daily partition with a decorator such as events$20261004, the documented window is 31 days in the past to 16 days in the future.
The gRPC Storage Write API and its stream types
The gRPC API replaces individual requests with write streams. A stream is a server-side object attached to a table; your client opens a long-lived gRPC connection and sends AppendRows requests on it, each carrying a batch of serialized protobuf rows. There are four stream types, and choosing one is the main design decision:
- Default stream. Exists for every table and needs no creation. Rows are visible as soon as the append is acknowledged. Semantics are at-least-once: if you retry an append whose response you never saw, the rows may land twice. This is the right choice for high-volume telemetry where occasional duplicates are tolerable or removed downstream.
- Committed type. You create the stream; rows are visible as soon as they are written. Appends can carry an offset, and the server only performs the write if the offset equals the next append offset of the stream. That check is what makes retries safe.
- Pending type. Rows are buffered invisibly until you finalize the stream and commit it. One commit can cover several streams atomically, which makes pending streams a streaming-protocol replacement for load jobs.
- Buffered type. Rows become visible only when you flush up to a row offset. Google describes it as an advanced type and generally recommends it only through Apache Beam.
The Java client's JsonStreamWriter converts JSON objects to protobuf using the table schema, which is why most examples use it.
Exactly-once with committed streams
Exactly-once writes need two things: a retry that cannot double-apply, and a durable record of how far you got. Committed streams provide the first through offsets. The second is your job, because a stream offset only means something if you store it next to the position in your source that produced it.
// Java: one committed stream per source partition, offsets tracked by the caller
BigQueryWriteClient client = BigQueryWriteClient.create();
TableName parent = TableName.of("proj", "analytics", "events");
WriteStream stream = WriteStream.newBuilder().setType(WriteStream.Type.COMMITTED).build();
WriteStream created = client.createWriteStream(
CreateWriteStreamRequest.newBuilder()
.setParent(parent.toString())
.setWriteStream(stream)
.build());
try (JsonStreamWriter writer =
JsonStreamWriter.newBuilder(created.getName(), created.getTableSchema()).build()) {
long offset = checkpoint.nextBigQueryOffset(); // 0 for a new stream
for (Batch b : source.batchesFrom(checkpoint.sourcePosition())) {
JSONArray rows = b.toJson();
ApiFuture<AppendRowsResponse> f = writer.append(rows, offset);
f.get(); // throws on failure; see retry rules below
offset += rows.length();
checkpoint.save(created.getName(), offset, b.endPosition());
}
}
client.finalizeWriteStream(created.getName()); // no more appends; rows are already visibleThe retry rules follow from the offset check. If an append fails with a transient status (Google's samples treat INTERNAL, CANCELLED and ABORTED as retryable), re-send the same rows with the same offset. If the server answers that the offset already exists, the earlier attempt succeeded and you should advance as if it had returned normally. If it answers that the offset is out of range, you skipped rows; that is a bug in your checkpoint logic, so stop rather than guess.
The checkpoint write is the subtle part. Save the stream name, the next offset and the source position together, atomically. After a crash you reopen the same stream and resume at the saved offset; any rows the dead process appended past that point are rejected as existing offsets. If you instead create a new stream on restart, the offset protection does not span streams and the replayed rows land twice.
Pending streams for atomic batches
A pending stream is the answer when a set of rows must appear all at once or not at all, for example a nightly reconciliation that replaces a day's data, or an export where a half-visible table would mislead dashboards. Write to one or more pending streams in parallel, finalize each with finalizeWriteStream, then call batchCommitWriteStreams with all their names. Until the commit, nothing is visible; after it, everything is. Compared with a load job you avoid staging files in Cloud Storage; compared with a committed stream you lose immediate visibility. A writer that crashes before the commit leaves nothing visible, and the whole unit must be rewritten.
Upserts and deletes with CDC
Plain streaming is append-only. The gRPC API also supports change data capture: you write rows with a _CHANGE_TYPE pseudocolumn whose value is UPSERT or DELETE, and BigQuery applies them against the table's primary key. The documented requirements are strict. The table must declare a primary key (composite keys up to 16 columns), writes must go through the default stream, and the format must be protobuf. An optional _CHANGE_SEQUENCE_NUMBER orders several changes to the same key. The table option max_staleness lets BigQuery apply changes in the background at least once within that interval, so queries can read a slightly stale but cheaply merged table.
CREATE TABLE analytics.customers (
customer_id STRING NOT NULL,
email STRING,
tier STRING,
updated_at TIMESTAMP,
PRIMARY KEY (customer_id) NOT ENFORCED
)
OPTIONS (max_staleness = INTERVAL 15 MINUTE);The limitations matter for design: a table with active CDC cannot take mutating DML such as UPDATE, DELETE or MERGE, and background merge jobs block copies, clones, snapshots and the Storage Read API. If you need those, stream changes into an append-only log table and MERGE into the target on a schedule instead.
Worked example: a clickstream pipeline
Consider a clickstream service emitting 20,000 events per second at about 1 KB each, roughly 20 MB/s or 1.7 TB per day, landing in a table partitioned by event date and clustered by user. Product wants dashboards under a minute behind real time; finance wants exact daily counts for billing. Those are different correctness requirements, and the design should not force one onto the other.
Events already flow through Pub/Sub. Pub/Sub delivery is at-least-once, so duplicates exist before BigQuery is involved. Making the BigQuery write exactly-once would not remove them; only a stable event ID can. So the design is:
- Producers stamp each event with an
event_idderived from its content and origin. - A writer pool (a Pub/Sub BigQuery subscription, or Dataflow, or your own service) appends to the default stream in batches of a few hundred rows to keep request sizes well inside the API limits. Acknowledge to Pub/Sub only after the append succeeds.
- Dashboards query the raw table directly and tolerate the occasional duplicate; measure the actual rate with a daily count of repeated
event_idvalues rather than assuming it. - Billing reads a deduplicated view, materialized nightly once the day has closed:
CREATE OR REPLACE TABLE billing.events_dedup_20261003 AS
SELECT * EXCEPT (rn) FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY ingest_time) AS rn
FROM analytics.events
WHERE event_date = '2026-10-03'
)
WHERE rn = 1;When would committed streams be worth their complexity? When the source has its own replayable offsets, such as Kafka partitions, and you cannot afford a deduplication pass, for instance because the table is the system of record. Then map each source partition to one committed stream, checkpoint as shown earlier, and accept owning the recovery code.
Monitoring and operations
BigQuery exposes ingestion history in INFORMATION_SCHEMA. The WRITE_API_TIMELINE views cover the gRPC API and the STREAMING_TIMELINE views cover the REST method, aggregated into one-minute intervals with byte totals broken down by error code (OK for successful requests). Alert on non-OK error codes and on bytes written falling below the expected rate; the second catches a writer that is healthy but disconnected from its source.
SELECT start_timestamp, error_code, SUM(total_input_bytes) AS bytes
FROM `region-us`.INFORMATION_SCHEMA.WRITE_API_TIMELINE_BY_PROJECT
WHERE start_timestamp > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR)
GROUP BY 1, 2
ORDER BY 1 DESC;Operationally, reuse connections. Creating a stream or connection per request exhausts connection quotas and adds latency; one long-lived writer per process (or per committed stream) is the pattern. Read the published request-size, throughput and connection quotas for your region before sizing a writer pool.
Failure modes
- Ignoring
insertErrors. The REST call returns success while some rows were rejected. Rows vanish with no exception anywhere. - A fresh
insertIdper attempt. Retries become duplicates because the IDs never match. - Retrying the default stream as if it were exactly-once. An ambiguous timeout followed by a retry can write the batch twice. Use committed streams or downstream deduplication.
- Schema drift. A producer adds a field; with
ignoreUnknownValuesit is silently dropped, and without it every row fails. On the gRPC path a writer built with the old schema rejects or drops the new field until it refreshes its schema. Add columns to the table first, then ship producers.
Trade-offs
| Need | Use | Cost of the choice |
|---|---|---|
| Minutes of latency acceptable, large files | Load jobs | Job quotas, staging files |
| Existing code, small volume | REST method with stable insertId | Best-effort dedup only |
| High volume, duplicates tolerable | Default stream | Read-side dedup for exact numbers |
| Exactly-once from a replayable source | Committed streams + checkpointed offsets | Stream lifecycle code |
| All-or-nothing batch | Pending streams + batch commit | No visibility until commit |
| Upserts and deletes by key | Default stream with CDC | No mutating DML, merge staleness |
What to do next
- Inventory every writer and classify it: REST method, default stream, committed, pending or CDC.
- For each REST writer, confirm it reads
insertErrorsand thatinsertIdis derived from the event, not generated per attempt. - Decide per table whether duplicates are tolerable; if not, either add read-side deduplication or move the writer to committed streams with checkpointed offsets.
- Partition streamed tables on an explicit event-time column rather than ingestion time.
- Add an INFORMATION_SCHEMA timeline query to monitoring with alerts on error codes and on throughput dropping below baseline.
- Look up current pricing and quotas for your region, then size batch sizes and writer counts from them.