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_id for streaming de-duplication.
  • Include _ingestion_time so 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_count vs expected): Drop in upload count signals streaming failures.
  • Pipeline freshness: Query SELECT MAX(_ingestion_time) FROM analytics.user_events and 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:

  1. Open Looker Studio and click "Create" → "Data Source"
  2. Select BigQuery as your connector
  3. Choose your project, dataset (analytics), and table (user_events)
  4. Enable "Enable Date Range Parameters" for time-filter performance
  5. 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:

  1. Application → Pub/Sub: Your services publish structured JSON events to a Pub/Sub topic with a defined Avro or Protobuf schema.
  2. Pub/Sub → Dataflow: A Dataflow streaming job subscribes to the topic, enriches events (geolocation lookup, session stitching, etc.), and handles dead letters.
  3. Dataflow → BigQuery: Enriched events are written to a partitioned, clustered BigQuery table via the Storage Write API.
  4. Materialized views: Hourly and daily aggregate views reduce query cost for dashboard workloads.
  5. 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.