From 71fede64b3afff4edd7ad285dacbe97d08a860d9 Mon Sep 17 00:00:00 2001 From: Vishal Gupta Date: Fri, 17 Jul 2026 05:01:12 +0000 Subject: [PATCH 1/2] Revert temp spanner table names --- pipeline/workflow/spanner-ingestion-workflow.yaml | 3 --- 1 file changed, 3 deletions(-) diff --git a/pipeline/workflow/spanner-ingestion-workflow.yaml b/pipeline/workflow/spanner-ingestion-workflow.yaml index 7a3ecd11..9c9b4a23 100644 --- a/pipeline/workflow/spanner-ingestion-workflow.yaml +++ b/pipeline/workflow/spanner-ingestion-workflow.yaml @@ -219,9 +219,6 @@ run_dataflow_job: projectId: '${spanner_project_id}' spannerInstanceId: '${spanner_instance_id}' spannerDatabaseId: '${spanner_database_id}' - # Revert this once switched to the new DB. - spannerObservationTableName: 'Observation_final_v3' - spannerTimeSeriesTableName: 'TimeSeries_final_v3' environment: numWorkers: ${numWorkers} machineType: ${machineType} From e11b2bb56753931fa50bccc37d6ebe80c1c20c95 Mon Sep 17 00:00:00 2001 From: Vishal Gupta Date: Mon, 27 Jul 2026 06:48:45 +0000 Subject: [PATCH 2/2] Remove dual writes for spanner DB --- pipeline/terraform/main.tf | 14 -------------- pipeline/workflow/aggregation-helper/main.py | 4 ++-- pipeline/workflow/cloudbuild.yaml | 3 +-- pipeline/workflow/deploy-services.yaml | 11 +++++------ .../workflow/ingestion-helper/clients/spanner.py | 12 ------------ pipeline/workflow/ingestion-helper/config.py | 1 - pipeline/workflow/ingestion-helper/dependencies.py | 1 - 7 files changed, 8 insertions(+), 38 deletions(-) diff --git a/pipeline/terraform/main.tf b/pipeline/terraform/main.tf index 2e399c65..ea0c851e 100644 --- a/pipeline/terraform/main.tf +++ b/pipeline/terraform/main.tf @@ -55,12 +55,6 @@ variable "spanner_database_id" { default = "dc-import-db" } -variable "spanner_graph_database_id" { - description = "Spanner Graph Database ID" - type = string - default = "dc-import-db" -} - variable "bq_dataset_id" { description = "BigQuery Dataset ID for aggregation" type = string @@ -193,10 +187,6 @@ resource "google_cloud_run_v2_service" "ingestion_helper" { name = "SPANNER_DATABASE_ID" value = var.spanner_database_id } - env { - name = "SPANNER_GRAPH_DATABASE_ID" - value = var.spanner_graph_database_id - } env { name = "GCS_BUCKET_ID" value = google_storage_bucket.import_bucket.name @@ -269,10 +259,6 @@ resource "google_cloud_run_v2_job" "aggregation_helper" { name = "SPANNER_DATABASE_ID" value = var.spanner_database_id } - env { - name = "SPANNER_GRAPH_DATABASE_ID" - value = var.spanner_graph_database_id - } env { name = "GCS_BUCKET_ID" value = google_storage_bucket.import_bucket.name diff --git a/pipeline/workflow/aggregation-helper/main.py b/pipeline/workflow/aggregation-helper/main.py index 84cf18db..cbd8b4e8 100644 --- a/pipeline/workflow/aggregation-helper/main.py +++ b/pipeline/workflow/aggregation-helper/main.py @@ -45,13 +45,13 @@ def create_orchestrator_config( connection_id = env.get("BQ_SPANNER_CONN_ID") project_id = env.get("PROJECT_ID") instance_id = env.get("SPANNER_INSTANCE_ID") - database_id = env.get("SPANNER_GRAPH_DATABASE_ID") + database_id = env.get("SPANNER_DATABASE_ID") location = env.get("LOCATION") if not connection_id or not project_id or not instance_id or not database_id: raise ValueError( f"Missing required environment variables. connection_id={connection_id}, " - f"project_id={project_id}, instance_id={instance_id}, database_id (SPANNER_GRAPH_DATABASE_ID)={database_id}" + f"project_id={project_id}, instance_id={instance_id}, database_id (SPANNER_DATABASE_ID)={database_id}" ) config_path = args.config_path or env.get("CONFIG_PATH") diff --git a/pipeline/workflow/cloudbuild.yaml b/pipeline/workflow/cloudbuild.yaml index 81a48961..7cde476c 100644 --- a/pipeline/workflow/cloudbuild.yaml +++ b/pipeline/workflow/cloudbuild.yaml @@ -21,7 +21,6 @@ substitutions: _SPANNER_PROJECT_ID: 'datcom-ci' _SPANNER_INSTANCE_ID: 'datcom-spanner-test' _SPANNER_DATABASE_ID: 'dc-test-db' - _SPANNER_GRAPH_DATABASE_ID: 'dc-test-db' _GCS_BUCKET_ID: 'datcom-ci-test' _GCS_MOUNT_BUCKET: 'datcom-ci-test' _BQ_DATASET_ID: 'spanner_model_dataset' @@ -56,7 +55,7 @@ steps: - '.' - '--config=deploy-services.yaml' - '--project=${_PROJECT_ID}' - - '--substitutions=_PROJECT_ID=${_PROJECT_ID},_SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},_SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},_SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},_SPANNER_GRAPH_DATABASE_ID=${_SPANNER_GRAPH_DATABASE_ID},_GCS_BUCKET_ID=${_GCS_BUCKET_ID},_LOCATION=${_LOCATION},_GCS_MOUNT_BUCKET=${_GCS_MOUNT_BUCKET},_BQ_DATASET_ID=${_BQ_DATASET_ID},_PROJECT_NUMBER=${_PROJECT_NUMBER},_BQ_SPANNER_CONN_ID=${_BQ_SPANNER_CONN_ID},_ENABLE_EMBEDDINGS=${_ENABLE_EMBEDDINGS},_VERSION=${_VERSION}' + - '--substitutions=_PROJECT_ID=${_PROJECT_ID},_SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},_SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},_SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},_GCS_BUCKET_ID=${_GCS_BUCKET_ID},_LOCATION=${_LOCATION},_GCS_MOUNT_BUCKET=${_GCS_MOUNT_BUCKET},_BQ_DATASET_ID=${_BQ_DATASET_ID},_PROJECT_NUMBER=${_PROJECT_NUMBER},_BQ_SPANNER_CONN_ID=${_BQ_SPANNER_CONN_ID},_ENABLE_EMBEDDINGS=${_ENABLE_EMBEDDINGS},_VERSION=${_VERSION}' dir: 'pipeline/workflow' # 3. Run E2E Tests on Staging diff --git a/pipeline/workflow/deploy-services.yaml b/pipeline/workflow/deploy-services.yaml index 6f1b3744..5f2f5864 100644 --- a/pipeline/workflow/deploy-services.yaml +++ b/pipeline/workflow/deploy-services.yaml @@ -19,9 +19,8 @@ substitutions: # Production config. _PROJECT_ID: 'datcom-import-automation-prod' _SPANNER_PROJECT_ID: 'datcom-store' - _SPANNER_INSTANCE_ID: 'dc-kg-test' - _SPANNER_DATABASE_ID: 'dc_graph_import' - _SPANNER_GRAPH_DATABASE_ID: 'dc_graph_2026_01_27' + _SPANNER_INSTANCE_ID: 'dc-graph-staging' + _SPANNER_DATABASE_ID: 'dc_graph_1' _GCS_BUCKET_ID: 'datcom-prod-imports' _LOCATION: 'us-central1' _GCS_MOUNT_BUCKET: 'datcom-volume-mount' @@ -36,7 +35,7 @@ substitutions: steps: - id: 'ingestion-helper-service' name: 'gcr.io/cloud-builders/gcloud' - args: ['run', 'deploy', 'ingestion-helper-service', '--image', '${_AR_REPO_URL}/datacommons-ingestion-helper:${_VERSION}', '--region', '${_LOCATION}', '--project', '${_PROJECT_ID}', '--no-allow-unauthenticated', '--timeout', '60m', '--set-env-vars', 'PROJECT_ID=${_PROJECT_ID},LOCATION=${_LOCATION},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},SPANNER_GRAPH_DATABASE_ID=${_SPANNER_GRAPH_DATABASE_ID},GCS_BUCKET_ID=${_GCS_BUCKET_ID}'] + args: ['run', 'deploy', 'ingestion-helper-service', '--image', '${_AR_REPO_URL}/datacommons-ingestion-helper:${_VERSION}', '--region', '${_LOCATION}', '--project', '${_PROJECT_ID}', '--no-allow-unauthenticated', '--timeout', '60m', '--set-env-vars', 'PROJECT_ID=${_PROJECT_ID},LOCATION=${_LOCATION},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},GCS_BUCKET_ID=${_GCS_BUCKET_ID}'] - id: 'import-helper-service' name: 'gcr.io/cloud-builders/gcloud' @@ -44,7 +43,7 @@ steps: - id: 'aggregation-helper-job' name: 'gcr.io/cloud-builders/gcloud' - args: ['run', 'jobs', 'deploy', 'aggregation-helper-job', '--image', '${_AR_REPO_URL}/datacommons-aggregation-helper:${_VERSION}', '--region', '${_LOCATION}', '--project', '${_PROJECT_ID}', '--tasks', '1', '--max-retries', '0', '--task-timeout', '60m', '--set-env-vars', 'PROJECT_ID=${_PROJECT_ID},LOCATION=${_LOCATION},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},SPANNER_GRAPH_DATABASE_ID=${_SPANNER_GRAPH_DATABASE_ID},GCS_BUCKET_ID=${_GCS_BUCKET_ID},BQ_DATASET_ID=${_BQ_DATASET_ID},BQ_SPANNER_CONN_ID=${_BQ_SPANNER_CONN_ID},ENABLE_EMBEDDINGS=${_ENABLE_EMBEDDINGS}'] + args: ['run', 'jobs', 'deploy', 'aggregation-helper-job', '--image', '${_AR_REPO_URL}/datacommons-aggregation-helper:${_VERSION}', '--region', '${_LOCATION}', '--project', '${_PROJECT_ID}', '--tasks', '1', '--max-retries', '0', '--task-timeout', '60m', '--set-env-vars', 'PROJECT_ID=${_PROJECT_ID},LOCATION=${_LOCATION},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},GCS_BUCKET_ID=${_GCS_BUCKET_ID},BQ_DATASET_ID=${_BQ_DATASET_ID},BQ_SPANNER_CONN_ID=${_BQ_SPANNER_CONN_ID},ENABLE_EMBEDDINGS=${_ENABLE_EMBEDDINGS}'] - id: 'import-automation-workflow' name: 'gcr.io/cloud-builders/gcloud' @@ -52,7 +51,7 @@ steps: - id: 'spanner-ingestion-workflow' name: 'gcr.io/cloud-builders/gcloud' - args: ['workflows', 'deploy', 'spanner-ingestion-workflow', '--project', '${_PROJECT_ID}', '--location', '${_LOCATION}', '--source', 'spanner-ingestion-workflow.yaml', '--set-env-vars', 'LOCATION=${_LOCATION},PROJECT_ID=${_PROJECT_ID},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_GRAPH_DATABASE_ID},PROJECT_NUMBER=${_PROJECT_NUMBER},DATAFLOW_TEMPLATE_PATH=${_DATAFLOW_TEMPLATE_PATH}ingestion-${_VERSION}.json'] + args: ['workflows', 'deploy', 'spanner-ingestion-workflow', '--project', '${_PROJECT_ID}', '--location', '${_LOCATION}', '--source', 'spanner-ingestion-workflow.yaml', '--set-env-vars', 'LOCATION=${_LOCATION},PROJECT_ID=${_PROJECT_ID},SPANNER_PROJECT_ID=${_SPANNER_PROJECT_ID},SPANNER_INSTANCE_ID=${_SPANNER_INSTANCE_ID},SPANNER_DATABASE_ID=${_SPANNER_DATABASE_ID},PROJECT_NUMBER=${_PROJECT_NUMBER},DATAFLOW_TEMPLATE_PATH=${_DATAFLOW_TEMPLATE_PATH}ingestion-${_VERSION}.json'] options: logging: CLOUD_LOGGING_ONLY diff --git a/pipeline/workflow/ingestion-helper/clients/spanner.py b/pipeline/workflow/ingestion-helper/clients/spanner.py index 92a39d4e..593b581d 100644 --- a/pipeline/workflow/ingestion-helper/clients/spanner.py +++ b/pipeline/workflow/ingestion-helper/clients/spanner.py @@ -56,7 +56,6 @@ def __init__(self, project_id: str, instance_id: str, database_id: str, - graph_database_id: str = None, location: str = None, models: list[dict] = None, embedding_space: int = 768, @@ -89,12 +88,6 @@ def __init__(self, database = instance.database(database_id) logging.info(f"Successfully initialized database: {database.name}") self.database = database - self.graph_database = database - if graph_database_id: - self.graph_database = instance.database(graph_database_id) - logging.info( - f"Successfully initialized graph database: {self.graph_database.name}" - ) self.project_id = project_id self.location = location self.embedding_space = embedding_space @@ -383,9 +376,6 @@ def _update(transaction: Transaction): try: self.database.run_in_transaction(_update) - # TODO: remove dual writes after switching to the prod setup. - if self.graph_database and self.graph_database.name != self.database.name: - self.graph_database.run_in_transaction(_update) logging.info( f"Updated IngestionHistory table for workflow {workflow_id}") except Exception as e: @@ -711,8 +701,6 @@ def _seed(transaction: Transaction): try: self.database.run_in_transaction(_seed) - if self.graph_database and self.graph_database.name != self.database.name: - self.graph_database.run_in_transaction(_seed) logging.info("Database seeded successfully.") except Exception as e: logging.error(f"Error seeding database: {e}") diff --git a/pipeline/workflow/ingestion-helper/config.py b/pipeline/workflow/ingestion-helper/config.py index 25455520..2ab13f4a 100644 --- a/pipeline/workflow/ingestion-helper/config.py +++ b/pipeline/workflow/ingestion-helper/config.py @@ -30,7 +30,6 @@ class EmbeddingSpec(BaseModel): SPANNER_PROJECT_ID = os.environ.get('SPANNER_PROJECT_ID') SPANNER_INSTANCE_ID = os.environ.get('SPANNER_INSTANCE_ID') SPANNER_DATABASE_ID = os.environ.get('SPANNER_DATABASE_ID') -SPANNER_GRAPH_DATABASE_ID = os.environ.get('SPANNER_GRAPH_DATABASE_ID') GCS_BUCKET_ID = os.environ.get('GCS_BUCKET_ID') LOCATION = os.environ.get('LOCATION') or os.environ.get('REGION') ENABLE_EMBEDDINGS = os.environ.get('ENABLE_EMBEDDINGS', 'false').lower() == 'true' diff --git a/pipeline/workflow/ingestion-helper/dependencies.py b/pipeline/workflow/ingestion-helper/dependencies.py index 8e1d4a76..fe7cdc23 100644 --- a/pipeline/workflow/ingestion-helper/dependencies.py +++ b/pipeline/workflow/ingestion-helper/dependencies.py @@ -33,7 +33,6 @@ def get_spanner_client() -> SpannerClient: config.SPANNER_PROJECT_ID, config.SPANNER_INSTANCE_ID, config.SPANNER_DATABASE_ID, - graph_database_id=config.SPANNER_GRAPH_DATABASE_ID, location=config.LOCATION, models=config.EMBEDDING_MODELS, embedding_space=config.EMBEDDING_SPACE,