Real-Time Analytics with BigQuery: Streaming Data from Pub/Sub to Dashboards
BigQuery's streaming insert API and native Pub/Sub integration make real-time analytics accessible without a dedicated streaming infrastructure team. This guide walks through a production-grade pipeline from event ingestion to live Looker Studio dashboards.
Most analytics pipelines start with a batch job that runs at midnight. You load yesterday's data, build your reports by morning, and stakeholders look at numbers that are already 12-plus hours old. That works fine for monthly business reviews. It doesn't work when you need to know right now whether your payment service is processing errors at an elevated rate or whether an A/B test is moving your conversion metric.
BigQuery's streaming capabilities have matured significantly over the past few years. The combination of Pub/Sub for durable event ingestion, Dataflow for stream processing, and BigQuery's native streaming API means you can get sub-minute latency from event to queryable data — without the operational overhead of running Kafka, Flink, or a dedicated streaming infrastructure team.
This guide is for data engineers and analytics engineers who need to build reliable real-time pipelines on GCP. We'll cover the full stack: Pub/Sub topic setup, streaming insert APIs, Dataflow templates, de-duplication strategies, partitioned destination tables, and live Looker Studio connections. We'll also be honest about the trade-offs and costs involved.
Understanding the Real-Time BigQuery Architecture
Before writing any code, it helps to understand the components and how they fit together.
Pub/Sub is Google's fully managed message bus. Publishers send messages to topics; subscribers pull from subscriptions. Pub/Sub gives you at-least-once delivery with automatic replication across zones. It handles traffic spikes gracefully — your event producers don't need to know or care about your downstream processing speed.
BigQuery Streaming Insert API allows you to insert rows into BigQuery tables immediately, outside of the normal batch load process. Streamed rows are available for queries within a few seconds. The trade-off: streaming inserts cost more than batch loads ($0.01 per 200 MB vs free for batch), and rows written via streaming inserts aren't immediately available for DML operations like UPDATE or DELETE.
Dataflow is the managed Apache Beam service on GCP. For many pipelines you don't need to write custom Beam code — Google provides pre-built Dataflow templates for common patterns like Pub/Sub to BigQuery.
BigQuery Storage Write API is the newer, preferred alternative to the legacy streaming insert API. It offers higher throughput, lower cost, and exactly-once semantics. For new pipelines, prefer the Storage Write API.
When Real-Time BigQuery Makes Sense
Real-time BigQuery is a good fit for:
- Operational metrics dashboards where 1-5 minute latency is acceptable
- Fraud detection or anomaly monitoring that runs as SQL queries
- Product analytics where you want same-session data visible
- Event-driven reporting that replaces overnight batch jobs
It's probably not the right tool for:
- Sub-second latency requirements (use a TSDB like BigTable or Prometheus)
- Complex stateful stream processing (use Flink or Spark Streaming)
- High-cardinality stream joins that would be expensive in BigQuery
Setting Up the Pub/Sub Foundation
Start by creating your Pub/Sub topic and subscription. We'll use the gcloud CLI throughout this guide.
# Create the topic that your application will publish to
gcloud pubsub topics create user-events --project=my-project
# Create a BigQuery subscription (this does the Pub/Sub → BigQuery routing natively)
gcloud pubsub subscriptions create user-events-to-bq --topic=user-events --bigquery-table=my-project:analytics.user_events --use-topic-schema --project=my-project
The BigQuery subscription type is the simplest path when you don't need stream processing — Pub/Sub writes directly to BigQuery without any Dataflow job. This works well for straightforward event logging.
For the direct BigQuery subscription to work, you need to define a Pub/Sub schema that matches your BigQuery table structure:
# Define a schema for structured event messages
gcloud pubsub schemas create user-event-schema --type=AVRO --definition='{
"type": "record",
"name": "UserEvent",
"fields": [
{"name": "event_id", "type": "string"},
{"name": "user_id", "type": "string"},
{"name": "event_type", "type": "string"},
{"name": "timestamp", "type": {"type": "long", "logicalType": "timestamp-micros"}},
{"name": "properties", "type": "string"}
]
}'
# Attach schema to your topic
gcloud pubsub topics update user-events --schema=user-event-schema --message-encoding=JSON
Designing Your BigQuery Destination Table
Your destination table structure has major implications for both query performance and streaming cost. Get this right before you start sending data.
CREATE TABLE analytics.user_events (
event_id STRING NOT NULL,
user_id STRING,
event_type STRING,
event_timestamp TIMESTAMP,
properties JSON,
_insert_id STRING, -- Used for de-duplication
_ingestion_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP()
)
PARTITION BY DATE(event_timestamp)
CLUSTER BY event_type, user_id
OPTIONS (
require_partition_filter = false,
partition_expiration_days = 365
);
Key decisions here:
- Partition by event timestamp, not ingestion time. This makes time-range queries efficient and predictable.
- Cluster by event_type and user_id — these are your most common filter columns.
- Include
_insert_idfor streaming de-duplication. - Include
_ingestion_timeso you can track pipeline lag.
Streaming Data with the BigQuery Storage Write API
For production pipelines with higher volume or exactly-once requirements, use the Storage Write API directly from your application. Here's a Python example:
from google.cloud.bigquery_storage_v1 import BigQueryWriteClient
from google.cloud.bigquery_storage_v1.types import (
AppendRowsRequest,
ProtoRows,
WriteStream,
)
from google.protobuf import descriptor_pb2
import proto
# Define your message schema (must match BigQuery table)
class UserEvent(proto.Message):
event_id = proto.Field(proto.STRING, number=1)
user_id = proto.Field(proto.STRING, number=2)
event_type = proto.Field(proto.STRING, number=3)
event_timestamp_micros = proto.Field(proto.INT64, number=4)
properties = proto.Field(proto.STRING, number=5)
def create_streaming_writer(project_id: str, dataset_id: str, table_id: str):
client = BigQueryWriteClient()
parent = client.table_path(project_id, dataset_id, table_id)
# Create a DEFAULT stream for at-least-once semantics
# Use COMMITTED stream for exactly-once
write_stream = WriteStream()
write_stream.type_ = WriteStream.Type.COMMITTED
write_stream = client.create_write_stream(
parent=parent,
write_stream=write_stream
)
return client, write_stream
def append_events(client, write_stream, events: list[dict]):
proto_rows = ProtoRows()
for event in events:
row = UserEvent(
event_id=event["event_id"],
user_id=event["user_id"],
event_type=event["event_type"],
event_timestamp_micros=int(event["timestamp"] * 1_000_000),
properties=str(event.get("properties", {})),
)
proto_rows.serialized_rows.append(row.serialize())
request = AppendRowsRequest()
request.write_stream = write_stream.name
request.proto_rows.rows.CopyFrom(proto_rows)
response = client.append_rows(iter([request]))
return response
Batching for Cost Efficiency
Streaming inserts are billed per byte inserted. Small, frequent inserts are expensive. Batch your events on the producer side before inserting:
import asyncio
from collections import defaultdict
from datetime import datetime
class StreamingBatcher:
def __init__(self, flush_interval_seconds=5, max_batch_size=1000):
self.buffer = []
self.flush_interval = flush_interval_seconds
self.max_batch_size = max_batch_size
self._last_flush = datetime.now()
async def add_event(self, event: dict):
self.buffer.append(event)
elapsed = (datetime.now() - self._last_flush).seconds
if len(self.buffer) >= self.max_batch_size or elapsed >= self.flush_interval:
await self.flush()
async def flush(self):
if not self.buffer:
return
batch = self.buffer.copy()
self.buffer.clear()
self._last_flush = datetime.now()
# Write batch to BigQuery
await self._write_to_bigquery(batch)
Using Dataflow for Stream Processing
When you need to transform, enrich, or aggregate events before loading to BigQuery, Dataflow is the right tool. The simplest path is the pre-built Pub/Sub to BigQuery Dataflow template:
gcloud dataflow jobs run pubsub-to-bigquery --gcs-location=gs://dataflow-templates/latest/PubSub_to_BigQuery --region=europe-west4 --staging-location=gs://my-project-dataflow/staging --parameters inputTopic=projects/my-project/topics/user-events,outputTableSpec=my-project:analytics.user_events,outputDeadletterTable=my-project:analytics.user_events_errors
Custom Dataflow Pipeline with Apache Beam
For enrichment or transformation logic, write a custom Beam pipeline:
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.io.gcp.pubsub import ReadFromPubSub
from apache_beam.io.gcp.bigquery import WriteToBigQuery, BigQueryDisposition
import json
from datetime import datetime
class EnrichEvent(beam.DoFn):
def process(self, message):
try:
event = json.loads(message.data.decode("utf-8"))
# Enrich with derived fields
event["date"] = datetime.utcfromtimestamp(
event["timestamp"]
).strftime("%Y-%m-%d")
event["hour"] = datetime.utcfromtimestamp(
event["timestamp"]
).hour
# Normalize event type
event["event_category"] = self._categorize_event(event["event_type"])
yield event
except Exception as e:
# Route malformed messages to dead letter
yield beam.pvalue.TaggedOutput("dead_letter", {
"raw_message": message.data.decode("utf-8"),
"error": str(e),
"ingestion_time": datetime.utcnow().isoformat(),
})
def _categorize_event(self, event_type: str) -> str:
if event_type.startswith("page_"):
return "navigation"
elif event_type.startswith("click_"):
return "interaction"
elif event_type.startswith("purchase_"):
return "conversion"
return "other"
def run_pipeline():
options = PipelineOptions()
options.view_as(StandardOptions).streaming = True
with beam.Pipeline(options=options) as p:
events = (
p
| "ReadPubSub" >> ReadFromPubSub(
topic="projects/my-project/topics/user-events",
with_attributes=True,
)
| "EnrichEvents" >> beam.ParDo(EnrichEvent()).with_outputs(
"dead_letter", main="valid_events"
)
)
# Write valid events to main table
events.valid_events | "WriteToBigQuery" >> WriteToBigQuery(
table="my-project:analytics.user_events",
write_disposition=BigQueryDisposition.WRITE_APPEND,
create_disposition=BigQueryDisposition.CREATE_NEVER,
)
# Write failed events to dead letter table
events.dead_letter | "WriteDeadLetter" >> WriteToBigQuery(
table="my-project:analytics.user_events_dead_letter",
write_disposition=BigQueryDisposition.WRITE_APPEND,
create_disposition=BigQueryDisposition.CREATE_IF_NEEDED,
)
if __name__ == "__main__":
run_pipeline()
Handling De-duplication
Pub/Sub guarantees at-least-once delivery. Under network partitions or retries, you may receive the same message more than once. BigQuery's streaming insert API has a built-in de-duplication mechanism using insertId, but it only works for 60 seconds after insert.
For a more robust approach, use a post-ingestion de-duplication query scheduled to run every 15 minutes:
-- Remove duplicates using MERGE on a temporary staging table
MERGE analytics.user_events AS target
USING (
SELECT *
FROM (
SELECT
*,
ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY _ingestion_time ASC) AS row_num
FROM analytics.user_events
WHERE DATE(event_timestamp) >= DATE_SUB(CURRENT_DATE(), INTERVAL 1 DAY)
)
WHERE row_num = 1
) AS source
ON target.event_id = source.event_id
AND DATE(target.event_timestamp) >= DATE_SUB(CURRENT_DATE(), INTERVAL 1 DAY)
WHEN NOT MATCHED THEN
INSERT ROW
WHEN MATCHED AND target._ingestion_time > source._ingestion_time THEN
DELETE;
Alternatively, use a streaming de-dup approach at the Dataflow layer with a sliding window and stateful processing.
Monitoring Pipeline Health
A streaming pipeline that silently falls behind is worse than one that fails loudly. Set up monitoring from day one.
# Create a Cloud Monitoring alert for Pub/Sub subscription backlog
gcloud alpha monitoring policies create --notification-channels=projects/my-project/notificationChannels/CHANNEL_ID --display-name="High Pub/Sub Backlog" --condition-display-name="subscription/num_undelivered_messages > 10000" --condition-filter='resource.type="pubsub_subscription" AND metric.type="pubsub.googleapis.com/subscription/num_undelivered_messages"' --condition-threshold-value=10000 --condition-threshold-comparison=COMPARISON_GT --condition-aggregations-per-series-aligner=ALIGN_MAX --condition-aggregations-alignment-period=300s
Key metrics to monitor:
- Pub/Sub subscription backlog (
subscription/num_undelivered_messages): Rising backlog means your consumer is falling behind. - Dataflow system lag (
dataflow.googleapis.com/job/system_lag): Time between oldest unprocessed element and current time. - BigQuery streaming errors (
bigquery.googleapis.com/storage/uploaded_row_countvs expected): Drop in upload count signals streaming failures. - Pipeline freshness: Query
SELECT MAX(_ingestion_time) FROM analytics.user_eventsand alert if it's more than 5 minutes old.
Building a Live Looker Studio Dashboard
With data streaming into BigQuery, connecting Looker Studio takes about 5 minutes:
- Open Looker Studio and click "Create" → "Data Source"
- Select BigQuery as your connector
- Choose your project, dataset (
analytics), and table (user_events) - Enable "Enable Date Range Parameters" for time-filter performance
- Add a calculated field for your key metric (e.g.,
COUNTIF(event_type = 'purchase'))
Optimizing Looker Studio Query Performance
Looker Studio translates your dashboard widgets into BigQuery SQL queries. Every time a user loads the dashboard or changes a date range, new queries run. Keep those queries fast:
-- Create a materialized view for dashboard aggregation
CREATE MATERIALIZED VIEW analytics.user_events_hourly AS
SELECT
DATE_TRUNC(event_timestamp, HOUR) AS event_hour,
event_type,
event_category,
COUNT(*) AS event_count,
COUNT(DISTINCT user_id) AS unique_users
FROM analytics.user_events
WHERE DATE(event_timestamp) >= DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)
GROUP BY 1, 2, 3;
Point your Looker Studio data source at analytics.user_events_hourly for most dashboard widgets. Reserve direct queries on user_events for drill-down scenarios where the user needs row-level detail.
Enabling Real-Time Refresh
Looker Studio doesn't auto-refresh by default. Enable it for operational dashboards:
- In report edit mode, click "File" → "Report Settings"
- Set "Data freshness" to "Real-time" (refreshes every 5 seconds when the tab is active)
- Note: real-time dashboards incur more BigQuery query cost — set a low cache TTL rather than 0
Cost Management for Streaming Pipelines
Streaming costs are easy to underestimate. Here's a rough breakdown for a pipeline processing 1 million events per day at 500 bytes each:
| Component | Daily Volume | Cost Estimate |
|---|---|---|
| Pub/Sub messages | 1M messages | ~$0.04 |
| Dataflow processing | 500 MB × some CPU | ~$0.50–$2.00 |
| BigQuery streaming inserts | 500 MB | ~$0.025 |
| BigQuery storage | 500 MB/day × 365 | ~$0.10/month |
| Looker Studio queries | Varies with users | $5/TB queried |
For most analytics pipelines at this scale, daily streaming cost lands under $5. The biggest variable is Dataflow autoscaling — set maxNumWorkers to cap spend during unexpected traffic spikes.
# Launch Dataflow with worker limits to cap cost
gcloud dataflow jobs run my-streaming-job --gcs-location=gs://dataflow-templates/latest/PubSub_to_BigQuery --region=europe-west4 --parameters inputTopic=projects/my-project/topics/user-events,outputTableSpec=my-project:analytics.user_events --max-workers=10 --num-workers=2
Putting It All Together
A production-ready real-time analytics pipeline on GCP looks like this:
- Application → Pub/Sub: Your services publish structured JSON events to a Pub/Sub topic with a defined Avro or Protobuf schema.
- Pub/Sub → Dataflow: A Dataflow streaming job subscribes to the topic, enriches events (geolocation lookup, session stitching, etc.), and handles dead letters.
- Dataflow → BigQuery: Enriched events are written to a partitioned, clustered BigQuery table via the Storage Write API.
- Materialized views: Hourly and daily aggregate views reduce query cost for dashboard workloads.
- Looker Studio: Connected to materialized views for dashboards, with direct table access for ad-hoc analysis.
The full pipeline from message published to dashboard updated typically runs in 30-90 seconds. For most operational analytics use cases, that's real enough.
For the next steps on BigQuery analytics, see our BigQuery ML guide to start building machine learning models directly on your streaming data. For broader GCP data architecture patterns, see our Dataflow guide on complex stream processing with Apache Beam.