Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 0 additions & 14 deletions pipeline/terraform/main.tf
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions pipeline/workflow/aggregation-helper/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
3 changes: 1 addition & 2 deletions pipeline/workflow/cloudbuild.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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
Expand Down
11 changes: 5 additions & 6 deletions pipeline/workflow/deploy-services.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -36,23 +35,23 @@ 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'
args: ['run', 'deploy', 'import-helper-service', '--image', '${_AR_REPO_URL}/datacommons-import-helper:${_VERSION}', '--region', '${_LOCATION}', '--project', '${_PROJECT_ID}', '--no-allow-unauthenticated', '--set-env-vars', 'PROJECT_ID=${_PROJECT_ID},LOCATION=${_LOCATION},PROJECT_NUMBER=${_PROJECT_NUMBER},GCS_BUCKET_ID=${_GCS_BUCKET_ID}']

- 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'
args: ['workflows', 'deploy', 'import-automation-workflow', '--project', '${_PROJECT_ID}', '--location', '${_LOCATION}', '--source', 'import-automation-workflow.yaml', '--set-env-vars', 'LOCATION=${_LOCATION},GCS_BUCKET_ID=${_GCS_BUCKET_ID},GCS_MOUNT_BUCKET=${_GCS_MOUNT_BUCKET},PROJECT_NUMBER=${_PROJECT_NUMBER}']

- 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
12 changes: 0 additions & 12 deletions pipeline/workflow/ingestion-helper/clients/spanner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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}")
Expand Down
1 change: 0 additions & 1 deletion pipeline/workflow/ingestion-helper/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
1 change: 0 additions & 1 deletion pipeline/workflow/ingestion-helper/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
3 changes: 0 additions & 3 deletions pipeline/workflow/spanner-ingestion-workflow.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand Down