Skip to content

Commit e282635

Browse files
authored
Revert "Add per-provenance counters to dataflow and spanner (#658)" (#694)
Reverting commit af33c21 to allow DCP rollout. Will resubmit afterwards.
1 parent 05497d0 commit e282635

5 files changed

Lines changed: 32 additions & 119 deletions

File tree

pipeline/ingestion/src/main/java/org/datacommons/ingestion/pipeline/GraphIngestionPipeline.java

Lines changed: 10 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,14 @@ public class GraphIngestionPipeline {
3333
private static final Logger LOGGER = LoggerFactory.getLogger(GraphIngestionPipeline.class);
3434
private static final Counter nodeInvalidTypeCounter =
3535
Metrics.counter(GraphIngestionPipeline.class, "mcf_nodes_without_type");
36+
private static final Counter nodeCounter =
37+
Metrics.counter(GraphIngestionPipeline.class, "graph_node_count");
38+
private static final Counter edgeCounter =
39+
Metrics.counter(GraphIngestionPipeline.class, "graph_edge_count");
40+
private static final Counter obsCounter =
41+
Metrics.counter(GraphIngestionPipeline.class, "graph_observation_count");
42+
private static final Counter timeSeriesCounter =
43+
Metrics.counter(GraphIngestionPipeline.class, "graph_timeseries_count");
3644

3745
// List of imports that require node combination.
3846
private static final List<String> IMPORTS_TO_COMBINE = List.of("Schema", "Place", "Provenance");
@@ -129,18 +137,6 @@ private static void processImport(
129137
boolean isBaseDc = options.getIsBaseDc();
130138
String provenance = ProvenanceUtils.getProvenanceDcid(importName, isBaseDc);
131139

132-
String shortName =
133-
importName.contains(":")
134-
? importName.substring(importName.lastIndexOf(':') + 1)
135-
: importName;
136-
137-
Counter nodeCounter = Metrics.counter(GraphIngestionPipeline.class, "node_count:" + shortName);
138-
Counter edgeCounter = Metrics.counter(GraphIngestionPipeline.class, "edge_count:" + shortName);
139-
Counter obsCounter =
140-
Metrics.counter(GraphIngestionPipeline.class, "observation_count:" + shortName);
141-
Counter timeSeriesCounter =
142-
Metrics.counter(GraphIngestionPipeline.class, "timeseries_count:" + shortName);
143-
144140
// 1. Prepare Deletes:
145141
// Generate mutations to delete existing data for this import/provenance.
146142
// Create a dummy signal if deletes are skipped, so downstream dependencies are satisfied
@@ -166,8 +162,7 @@ private static void processImport(
166162
PipelineUtils.InputFormat format = PipelineUtils.resolveFormat(graphPath);
167163
if (format == PipelineUtils.InputFormat.TFRECORD) {
168164
TfRecordProcessingResult result =
169-
processTfRecordImport(
170-
pipeline, importName, graphPath, isBaseDc, obsCounter, timeSeriesCounter);
165+
processTfRecordImport(pipeline, importName, graphPath, isBaseDc);
171166
writeToSpanner(
172167
pipeline,
173168
spannerClient,
@@ -331,12 +326,7 @@ private static class TfRecordProcessingResult {
331326
}
332327

333328
private static TfRecordProcessingResult processTfRecordImport(
334-
Pipeline pipeline,
335-
String importName,
336-
String graphPath,
337-
boolean isBaseDc,
338-
Counter obsCounter,
339-
Counter timeSeriesCounter) {
329+
Pipeline pipeline, String importName, String graphPath, boolean isBaseDc) {
340330
PCollection<McfOptimizedGraph> optGraph =
341331
PipelineUtils.readOptimizedMcfGraph(importName, graphPath, pipeline);
342332
PCollection<TimeSeries> uniqueSeries =

pipeline/workflow/ingestion-helper/clients/spanner.py

Lines changed: 11 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -401,34 +401,24 @@ def update_import_version_history(self,
401401
logging.info(
402402
f"Updating ImportVersionHistory table for workflow {workflow_id}")
403403

404-
m = metrics if metrics else {}
405-
import_metrics = m.get('import_metrics', {})
406-
407404
def _insert(transaction: Transaction):
408405
columns = [
409406
"ImportName", "Version", "UpdateTimestamp",
410407
"WorkflowExecutionID", "Status", "ExecutionTime",
411408
"NodeCount", "EdgeCount", "ObservationCount",
412409
"TimeSeriesCount", "Comment"
413410
]
411+
m = metrics if metrics else {}
414412
version_history_values = []
415413
for import_json in import_list_json:
416-
import_name = import_json.get('importName')
417-
if not import_name:
418-
continue
419-
420-
short_name = import_name.split(':')[-1]
421-
counts = import_metrics.get(short_name) or import_json
422-
423414
version_history_values.append([
424-
import_name, import_json.get('latestVersion'),
415+
import_json['importName'], import_json['latestVersion'],
425416
spanner.COMMIT_TIMESTAMP, workflow_id, status,
426417
m.get('execution_time'),
427-
counts.get('node_count') or counts.get('nodeCount'),
428-
counts.get('edge_count') or counts.get('edgeCount'),
429-
counts.get('obs_count') or counts.get('obsCount'),
430-
counts.get('ts_count') or counts.get('tsCount'),
431-
"ingestion-workflow:" + workflow_id
418+
m.get('node_count'),
419+
m.get('edge_count'),
420+
m.get('obs_count'),
421+
m.get('ts_count'), "ingestion-workflow:" + workflow_id
432422
])
433423

434424
if version_history_values:
@@ -522,27 +512,22 @@ def update_version_history(self,
522512
import_name = import_name.split(':')[-1]
523513
logging.info(f"Updating version history for {import_name} to {version}")
524514

525-
m = metrics if metrics else {}
526-
node_count = m.get('node_count')
527-
edge_count = m.get('edge_count')
528-
obs_count = m.get('obs_count')
529-
ts_count = m.get('ts_count')
530-
531515
def _record(transaction: Transaction):
532516
columns = [
533517
"ImportName", "Version", "UpdateTimestamp",
534518
"WorkflowExecutionID", "Status", "ExecutionTime",
535519
"NodeCount", "EdgeCount", "ObservationCount",
536520
"TimeSeriesCount", "Comment"
537521
]
522+
m = metrics if metrics else {}
538523
values = [[
539524
import_name, version, spanner.COMMIT_TIMESTAMP, workflow_id,
540525
status,
541526
m.get('execution_time'),
542-
node_count,
543-
edge_count,
544-
obs_count,
545-
ts_count, comment
527+
m.get('node_count'),
528+
m.get('edge_count'),
529+
m.get('obs_count'),
530+
m.get('ts_count'), comment
546531
]]
547532
transaction.insert(table="ImportVersionHistory",
548533
columns=columns,

pipeline/workflow/ingestion-helper/clients/spanner_test.py

Lines changed: 0 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -702,45 +702,7 @@ def run_in_transaction_side_effect(callback, *args, **kwargs):
702702
None, None, None, None, None, None, "ingestion-workflow:wf-789"
703703
]])
704704

705-
@patch('google.cloud.spanner.Client')
706-
def test_update_import_version_history_per_import_counts(self, mock_spanner_client):
707-
mock_instance = MagicMock()
708-
mock_db = MagicMock()
709-
mock_spanner_client.return_value.instance.return_value = mock_instance
710-
mock_instance.database.return_value = mock_db
711-
712-
mock_transaction = MagicMock()
713-
def run_in_transaction_side_effect(callback, *args, **kwargs):
714-
return callback(mock_transaction, *args, **kwargs)
715-
mock_db.run_in_transaction.side_effect = run_in_transaction_side_effect
716-
717-
client = SpannerClient("project", "instance", "database")
718-
719-
import_list = [
720-
{"importName": "import1", "latestVersion": "v1.0"},
721-
{"importName": "import2", "latestVersion": "v2.0"}
722-
]
723-
metrics = {
724-
'execution_time': 120,
725-
'import_metrics': {
726-
"import1": {'node_count': 10, 'edge_count': 20, 'ts_count': 5, 'obs_count': 15},
727-
"import2": {'node_count': 30, 'edge_count': 40, 'ts_count': 25, 'obs_count': 35},
728-
}
729-
}
730-
client.update_import_version_history(
731-
import_list, workflow_id="wf-123", status="SUCCESS", metrics=metrics
732-
)
733-
734-
mock_transaction.insert.assert_called_once()
735-
_, kwargs = mock_transaction.insert.call_args
736-
self.assertEqual(kwargs['table'], 'ImportVersionHistory')
737-
self.assertEqual(kwargs['values'], [
738-
["import1", "v1.0", spanner.COMMIT_TIMESTAMP, "wf-123", "SUCCESS", 120, 10, 20, 15, 5, "ingestion-workflow:wf-123"],
739-
["import2", "v2.0", spanner.COMMIT_TIMESTAMP, "wf-123", "SUCCESS", 120, 30, 40, 35, 25, "ingestion-workflow:wf-123"]
740-
])
741-
742705

743706
if __name__ == '__main__':
744707
unittest.main()
745708

746-

pipeline/workflow/ingestion-helper/routes/database.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ def acquire_ingestion_lock(req: LockAcquireRequest, spanner: SpannerClient = Dep
5858
status_ok = spanner.acquire_lock(req.workflowId, req.timeout)
5959
if not status_ok:
6060
raise HTTPException(
61-
status_code=503,
61+
status_code=400,
6262
detail=f"Failed to acquire lock: Lock already held or acquisition timed out for workflow {req.workflowId}"
6363
)
6464
return BaseResponse(status=ResponseStatus.OK)

pipeline/workflow/ingestion-helper/utils/imports.py

Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -127,8 +127,7 @@ def get_ingestion_metrics(project_id, location, job_id):
127127
job_id: The Dataflow job ID.
128128
129129
Returns:
130-
A dictionary containing 'obs_count', 'node_count', 'edge_count', 'ts_count', 'execution_time',
131-
and 'import_metrics'.
130+
A dictionary containing 'obs_count', 'node_count', 'edge_count', 'ts_count', and 'execution_time'.
132131
"""
133132
dataflow = build('dataflow', 'v1b3', cache_discovery=False)
134133
# Fetch Dataflow metrics
@@ -137,7 +136,6 @@ def get_ingestion_metrics(project_id, location, job_id):
137136
obs_count = 0
138137
ts_count = 0
139138
execution_time = 0
140-
import_metrics = {}
141139
if project_id and job_id:
142140
try:
143141
# Fetch Job details for execution time
@@ -161,35 +159,14 @@ def get_ingestion_metrics(project_id, location, job_id):
161159
jobId=job_id).execute()
162160
for metric in metrics.get('metrics', []):
163161
name = metric['name']['name']
164-
scalar = int(metric.get('scalar', 0))
165-
match = re.match(r"^([^:]+)(?::(.+))?$", name)
166-
metric_type = match.group(1) if match else name
167-
imp_name = match.group(2).split(':')[-1] if match and match.group(2) else None
168-
169-
if imp_name and imp_name not in import_metrics:
170-
import_metrics[imp_name] = {
171-
'node_count': 0,
172-
'edge_count': 0,
173-
'obs_count': 0,
174-
'ts_count': 0,
175-
}
176-
177-
if metric_type == 'node_count':
178-
node_count += scalar
179-
if imp_name:
180-
import_metrics[imp_name]['node_count'] = scalar
181-
elif metric_type == 'edge_count':
182-
edge_count += scalar
183-
if imp_name:
184-
import_metrics[imp_name]['edge_count'] = scalar
185-
elif metric_type == 'observation_count':
186-
obs_count += scalar
187-
if imp_name:
188-
import_metrics[imp_name]['obs_count'] = scalar
189-
elif metric_type == 'timeseries_count':
190-
ts_count += scalar
191-
if imp_name:
192-
import_metrics[imp_name]['ts_count'] = scalar
162+
if name == 'graph_node_count':
163+
node_count += int(metric['scalar'])
164+
elif name == 'graph_edge_count':
165+
edge_count += int(metric['scalar'])
166+
elif name == 'graph_observation_count':
167+
obs_count += int(metric['scalar'])
168+
elif name == 'graph_timeseries_count':
169+
ts_count += int(metric['scalar'])
193170
except HttpError as e:
194171
logging.error(
195172
f"Error fetching dataflow metrics for job {job_id}: {e}")
@@ -198,6 +175,5 @@ def get_ingestion_metrics(project_id, location, job_id):
198175
'node_count': node_count,
199176
'edge_count': edge_count,
200177
'ts_count': ts_count,
201-
'execution_time': execution_time,
202-
'import_metrics': import_metrics,
178+
'execution_time': execution_time
203179
}

0 commit comments

Comments
 (0)