This document describes how to use reverse extract, transform, and load (ETL) pipelines to move and continuously synchronize graph data from BigQuery to Spanner Graph. It covers the following key aspects:
- Common use cases for reverse ETL with graph data.
- The steps involved in a reverse ETL pipeline.
- Strategies for managing graph data changes, including insertions, updates, and deletions.
- Methods for orchestrating and maintaining reverse ETL pipelines.
- Best practices for optimizing your reverse ETL process.
To use reverse ETL to export data from BigQuery to Spanner, see Export data to Spanner.
BigQuery performs complex data manipulation at scale as an analytical processing platform, while Spanner is optimized for use cases that require high QPS and low serving latency. Spanner Graph and BigQuery integrate effectively to prepare graph data in BigQuery analytics pipelines, enabling Spanner to serve low-latency graph traversals.
Before you begin
Create a Spanner instance with a database that contains graph data. For more information, see Set up and query Spanner Graph.
In BigQuery, create an Enterprise or Enterprise Plus tier slot reservation. You can reduce BigQuery compute costs when you run exports to Spanner Graph. To do this, set a baseline slot capacity of zero and enable autoscaling.
Grant Identity and Access Management (IAM) roles that give users the necessary permissions to perform each task in this document.
Required roles
To get the permissions that you need to export BigQuery graph data to Spanner Graph, ask your administrator to grant you the following IAM roles on your project:
-
Export data from a BigQuery table:
BigQuery Data Viewer (
roles/bigquery.dataViewer) -
Run an export job:
BigQuery User (
roles/bigquery.user) -
View parameters of the Spanner instance:
Cloud Spanner Viewer (
roles/spanner.viewer) -
Write data to a Spanner Graph table:
Cloud Spanner Database User (
roles/spanner.databaseUser)
For more information about granting roles, see Manage access to projects, folders, and organizations.
You might also be able to get the required permissions through custom roles or other predefined roles.
Reverse ETL use cases
The following are example use cases. After you analyze and process data in BigQuery, you can move the data to Spanner Graph using reverse ETL.
Data aggregation and summarization - Use BigQuery to compute aggregates over granular data to make it more suitable for operational use cases.
Data transformation and enrichment - Use BigQuery to cleanse and standardize data received from different data sources.
Data filtering and selection - Use BigQuery to filter a large dataset for analytical purposes. For example, you might filter out data that is not required for real-time applications.
Feature preprocessing and engineering - In BigQuery, use the ML.TRANSFORM function to transform data, or the ML.FEATURE_CROSS function to create feature crosses of input features. Then, use reverse ETL to move the resulting data into Spanner Graph.
Understand the reverse ETL pipeline
Data moves from BigQuery to Spanner Graph in a reverse ETL pipeline in two steps:
BigQuery uses slots assigned to the pipeline job to extract and transform source data.
The BigQuery reverse ETL pipeline uses Spanner APIs to load data into a provisioned Spanner instance.
The following diagram shows the steps in a reverse ETL pipeline:
Figure 1. BigQuery reverse ETL pipeline process
Manage graph data changes
You can use reverse ETL to do the following:
Load a graph dataset from BigQuery to Spanner Graph.
Synchronize Spanner Graph data with ongoing updates from a dataset in BigQuery.
You configure a reverse ETL pipeline with a SQL query to specify the source data
and the transformation to apply. The pipeline loads all data that satisfies the
WHERE clause of the SELECT statement into Spanner using an
upsert operation. An upsert operation is equivalent to
INSERT OR UPDATE
statements. It inserts new rows and updates existing rows in tables that store
graph data. The pipeline bases new and updated rows on a Spanner
table primary key.
Insert and update data for tables with load order dependencies
Spanner Graph schema design best practices recommend using interleaved tables and foreign keys. If you use interleaved tables or enforced foreign keys, you must load node and edge data in a specific order. This ensures referenced rows exist before you create the referencing row. For more information, see Create interleaved tables.
The following example graph input table schema uses an interleaved table and a foreign key constraint to model the relationship between a person and their accounts:
CREATE TABLE Person (
id INT64 NOT NULL,
name STRING(MAX)
) PRIMARY KEY (id);
CREATE TABLE Account (
id INT64 NOT NULL,
create_time TIMESTAMP,
is_blocked BOOL,
type STRING(MAX)
) PRIMARY KEY (id);
CREATE TABLE PersonOwnAccount (
id INT64 NOT NULL,
account_id INT64 NOT NULL,
create_time TIMESTAMP,
CONSTRAINT FK_Account FOREIGN KEY (account_id) REFERENCES Account (id)
) PRIMARY KEY (id, account_id),
INTERLEAVE IN PARENT Person ON DELETE CASCADE;
CREATE PROPERTY GRAPH FinGraph
NODE TABLES (
Person,
Account
)
EDGE TABLES (
PersonOwnAccount
SOURCE KEY (id) REFERENCES Person
DESTINATION KEY (account_id) REFERENCES Account
LABEL Owns
);
In this example schema, PersonOwnAccount is an interleaved table in Person.
Load elements in the Person table before elements in the
PersonOwnAccount table. Additionally, the foreign key constraint on
PersonOwnAccount ensures a matching row exists in Account, the edge
relationship target. Therefore, load the Account table before the
PersonOwnAccount table. The following list summarizes this schema's load order
dependencies:
Follow these steps to load the data:
- Load
PersonbeforePersonOwnAccount. - Load
AccountbeforePersonOwnAccount.
Spanner enforces the referential integrity constraints in the
example schema. If the pipeline attempts to create a row in the
PersonOwnAccount table without a matching row in either the Person table or
Account table, Spanner returns an error. The pipeline then
fails.
This example reverse ETL pipeline uses
EXPORTDATA
statements in BigQuery to export data from the Person,
Account, and PersonOwnAccount tables in a dataset to meet load order
dependencies:
BEGIN
EXPORT DATA OPTIONS (
uri="https://spanner.googleapis.com/projects/PROJECT_ID/instances/INSTANCE_ID/databases/DATABASE_ID",
format='CLOUD_SPANNER',
spanner_options="""{
"table": "Person",
"priority": "LOW",
"tag" : "graph_data_load_person"
}"""
) AS
SELECT
id,
name
FROM
DATASET_NAME.Person;
EXPORT DATA OPTIONS (
uri="https://spanner.googleapis.com/projects/PROJECT_ID/instances/INSTANCE_ID/databases/DATABASE_ID",
format='CLOUD_SPANNER',
spanner_options="""{
"table": "Account",
"priority": "LOW",
"tag" : "graph_data_load_account"
}"""
) AS
SELECT
id,
create_time,
is_blocked,
type
FROM
DATASET_NAME.Account;
EXPORT DATA OPTIONS (
uri="https://spanner.googleapis.com/projects/PROJECT_ID/instances/INSTANCE_ID/databases/DATABASE_ID",
format='CLOUD_SPANNER',
spanner_options="""{
"table": "PersonOwnAccount",
"priority": "LOW",
"tag" : "graph_data_load_person_own_account"
}"""
) AS
SELECT
id,
account_id,
create_time
FROM
DATASET_NAME.PersonOwnAccount;
END;
Synchronize data
To synchronize BigQuery with Spanner Graph, use reverse ETL pipelines. You can configure a pipeline to do one of the following:
Apply any insertions and updates from the BigQuery source to the Spanner Graph target table. You can add schema elements to the target tables to logically communicate deletes and remove target table rows on a schedule.
Use a time series function that applies insert and update operations and identifies delete operations.
Referential integrity constraints
Unlike Spanner, BigQuery doesn't enforce primary and foreign key constraints. If your BigQuery data doesn't conform to the constraints you create on your Spanner tables, the reverse ETL pipeline might fail when loading that data.
Reverse ETL automatically groups data into batches that don't exceed the maximum mutation per commit limit and atomically applies the batches to a Spanner table in an arbitrary order. If a batch contains data that fails a referential integrity check, Spanner doesn't load that batch. Examples of such failures include an interleaved child row lacking a parent row or an enforced foreign key column without a matching value in the referenced column. If a batch fails a check, the pipeline fails with an error, and the pipeline stops loading batches.
Understand referential integrity constraint errors
The following examples show referential integrity constraint errors that you might encounter:
Resolve foreign key constraint errors
Error: "Foreign key constraint
FK_Accountis violated on tablePersonOwnAccount. Can't find referenced values inAccount(id)"Cause: A row insert into the
PersonOwnAccounttable failed because a matching row in theAccounttable, which theFK_Accountforeign key requires, is missing.
Resolve parent row missing errors
Error: "Parent row for row [15,1] in table
PersonOwnAccountis missing"Cause: A row insert into
PersonOwnAccount(id: 15andaccount_id: 1) failed because a parent row in thePersontable (id: 15) is missing.
To reduce the risk of referential integrity errors, consider the following options. Each option has tradeoffs.
- Relax the constraints to allow Spanner Graph to load data.
- Add logic to your pipeline to omit rows that violate referential integrity constraints.
Relax referential integrity
One option to avoid referential integrity errors when loading data is to relax the constraints so that Spanner doesn't enforce referential integrity.
You can create interleaved tables with the
INTERLEAVE INclause to use the same physical row interleaving characteristics. If you useINTERLEAVE INinstead ofINTERLEAVE IN PARENT, Spanner doesn't enforce referential integrity, though queries benefit from the co-location of related tables.You can create informational foreign keys by using the
NOT ENFORCEDoption. TheNOT ENFORCEDoption provides query optimization benefits. Spanner doesn't, however, enforce referential integrity.
For example, to create the edge input table without referential integrity checks, you can use this DDL:
CREATE TABLE PersonOwnAccount (
id INT64 NOT NULL,
account_id INT64 NOT NULL,
create_time TIMESTAMP,
CONSTRAINT FK_Account FOREIGN KEY (account_id) REFERENCES Account (id) NOT ENFORCED
) PRIMARY KEY (id, account_id),
INTERLEAVE IN Person;
Respect referential integrity in reverse ETL pipelines
To ensure the pipeline loads only rows that satisfy the referential integrity
checks, include only PersonOwnAccount rows that have matching rows in the
Person and Account tables. Then, preserve the load order, so
Spanner loads Person and Account rows before the
PersonOwnAccount rows that refer to them.
EXPORT DATA OPTIONS (
uri="https://spanner.googleapis.com/projects/PROJECT_ID/instances/INSTANCE_ID/databases/DATABASE_ID",
format='CLOUD_SPANNER',
spanner_options="""{
"table": "PersonOwnAccount",
"priority": "LOW",
"tag" : "graph_data_load_person_own_account"
}"""
) AS
SELECT
poa.id,
poa.account_id,
poa.create_time
FROM `PROJECT_ID.DATASET_NAME.PersonOwnAccount` poa
JOIN `PROJECT_ID.DATASET_NAME.Person` p ON (poa.id = p.id)
JOIN `PROJECT_ID.DATASET_NAME.Account` a ON (poa.account_id = a.id)
WHERE poa.id = p.id
AND poa.account_id = a.id;
Delete graph elements
Reverse ETL pipelines use upsert operations. Because upsert operations are
equivalent to
INSERT OR UPDATE
statements, a pipeline can only synchronize rows that exist in the source data
at runtime. This means the pipeline excludes deleted rows. If you delete data
from BigQuery, a reverse ETL pipeline can't directly remove the
same data from Spanner Graph.
You can use one of the following options to handle deletions from BigQuery source tables:
Perform a logical or soft delete in the source
To logically mark rows for deletion, use a deleted flag in BigQuery. Then create a column in the target Spanner table to which you can propagate the flag. When the reverse ETL applies the pipeline updates, delete rows that have this flag in Spanner. You can find and delete such rows explicitly using partitioned DML. Alternatively, implicitly delete rows by configuring a TTL (time to live) column with a date that depends on the delete flag column. Write Spanner queries to exclude these logically deleted rows. This ensures that Spanner excludes these rows from results before scheduled deletion. After the reverse ETL pipeline runs to completion, Spanner reflects the logical deletes in its rows. You can then delete rows from BigQuery.
This example adds an is_deleted column to the PersonOwnAccount table in
Spanner. It then adds an expired_ts_generated column that
depends on the is_deleted value. The TTL policy schedules affected rows for
deletion because the date in the generated column is earlier than the
DELETION POLICY threshold.
ALTER TABLE PersonOwnAccount
ADD COLUMN is_deleted BOOL DEFAULT (FALSE);
ALTER TABLE PersonOwnAccount ADD COLUMN
expired_ts_generated TIMESTAMP AS (IF(is_deleted,
TIMESTAMP("1970-01-01 00:00:00+00"),
TIMESTAMP("9999-01-01 00:00:00+00"))) STORED HIDDEN;
ALTER TABLE PersonOwnAccount
ADD ROW DELETION POLICY (OLDER_THAN(expired_ts_generated, INTERVAL 0 DAY));
Use BigQuery change history for INSERT, UPDATE and logical deletes
You can track changes to a BigQuery table using its change
history. Use the GoogleSQL
CHANGES
function to find rows that changed within a specific time interval. Then, use
the deleted row information with a reverse ETL pipeline. You can set up the
pipeline to set an indicator, like a deleted flag or expiration date, in the
Spanner table. This indicator marks rows for deletion in the
Spanner tables.
Use the results from the CHANGES time series function to decide which rows
from the source table to include in the load of your reverse ETL pipeline.
The pipeline includes rows with _CHANGE_TYPE as INSERT or UPDATE as
upserts if the row exists in the source table. The current row from the source
table provides the most recent data.
Use rows with _CHANGE_TYPE as DELETE that do not have existing rows in
the source table to set an indicator in the Spanner table, such
as a deleted flag or row expiration date.
Your export query must account for the order of insertions and deletions in BigQuery. For example, consider a row deleted at time T1 and a new row inserted at a later time T2. If both map to the same Spanner table row, the export must preserve the effects of these events in their original order.
If set, the delete indicator marks rows for deletion in the Spanner tables.
For example, you could add a column to a Spanner input table to store each row's expiration date. Then, create a deletion policy that uses these expiration dates.
The following example shows how to add a column to store the expiration dates of the table's rows.
ALTER TABLE PersonOwnAccount ADD COLUMN expired_ts TIMESTAMP;
ALTER TABLE PersonOwnAccount
ADD ROW DELETION POLICY (OLDER_THAN(expired_ts, INTERVAL 1 DAY));
To use the CHANGES function on a table in BigQuery, set the
table's
enable_change_history option
to TRUE:
ALTER TABLE `PROJECT_ID.DATASET_NAME.PersonOwnAccount`
SET OPTIONS (enable_change_history=TRUE);
The following example shows how you can use reverse ETL to update new or changed
rows and set the expiration date for rows marked for deletion. A left join with
the PersonOwnAccount table gives the query information about each row's
current status.
EXPORT DATA OPTIONS (
uri="https://spanner.googleapis.com/projects/PROJECT_ID/instances/INSTANCE_ID/databases/DATABASE_ID",
format='CLOUD_SPANNER',
spanner_options="""{
"table": "PersonOwnAccount",
"priority": "LOW",
"tag" : "graph_data_delete_via_reverse_etl"
}"""
) AS
SELECT
DISTINCT
IF (changes._CHANGE_TYPE = 'DELETE', changes.id, poa.id) AS id,
IF (changes._CHANGE_TYPE = 'DELETE', changes.account_id, poa.account_id) AS account_id,
IF (changes._CHANGE_TYPE = 'DELETE', changes.create_time, poa.create_time) AS create_time,
IF (changes._CHANGE_TYPE = 'DELETE', changes._CHANGE_TIMESTAMP, NULL) AS expired_ts
FROM
CHANGES(TABLE `PROJECT_ID.DATASET_NAME.PersonOwnAccount`,
TIMESTAMP_TRUNC(TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY), DAY),
TIMESTAMP_TRUNC(CURRENT_TIMESTAMP(), DAY)) changes
LEFT JOIN `PROJECT_ID.DATASET_NAME.PersonOwnAccount` poa
ON (poa.id = changes.id
AND poa.account_id = changes.account_id)
WHERE (changes._CHANGE_TYPE = 'DELETE'
AND poa.id IS NULL)
OR (changes._CHANGE_TYPE IN ( 'UPDATE', 'INSERT')
AND poa.id IS NOT NULL );
The example query uses a LEFT JOIN with the source table to preserve order.
This join ensures that DELETE change records are ignored for rows deleted and
then recreated within the query change history interval. The pipeline preserves
the valid, new row.
When you delete rows, the pipeline populates the expired_ts column in the
corresponding Spanner Graph row using the