Skip to content

CensusCountyBusinessPatterns: optimize pipeline runtime and update validation - #2219

Open
kartik-s21 wants to merge 61 commits into
datacommonsorg:masterfrom
kartik-s21:fix-census-cbp-validation
Open

kartik-s21 wants to merge 61 commits into
datacommonsorg:masterfrom
kartik-s21:fix-census-cbp-validation

Conversation

@kartik-s21

@kartik-s21 kartik-s21 commented Sep 10, 2026 •

Copy link
Copy Markdown
Contributor

Summary

This PR resolves validation failures, eliminates redundant network calls, and significantly improves runtime performance for CensusCountyBusinessPatterns.

Key Changes

  1. Pipeline & Concurrency Optimization:

    • Local Schema Caching: Download gs://unresolved_mcf/scripts/statvar/stat_vars.mcf once locally before shard processing, eliminating >13 GB of redundant concurrent GCS downloads across parallel workers.
    • Parallelized Sharding & Header Injection: Streamed sharding in Step 1 using split --filter to inject headers directly on write, avoiding expensive sed -i in-place disk rewrites.
    • Dynamic Concurrency: Increased shard chunk size to 1,000,000 rows and dynamically scaled PARALLELISM up to 48 workers on 64-core VMs.
    • I/O Streaming & Multiprocessing in main.py: Switched to streaming decompressed inputs directly into csv.DictReader and added ThreadPoolExecutor for zip downloads, eliminating multi-GB memory buffers and reducing runtime.
  2. Fault Tolerance, Cleanup & Error Propagation:

    • Working Directory Hygiene: Pre-cleans $SHARD_DIR/*_shard_*.csv and $OUTPUT_FINAL_DIR/output_* before Step 1 to prevent ingestion of stale artifacts from prior runs.
    • Fail-Fast Worker Termination: Background job failures in Step 1 or Step 2 immediately kill remaining active jobs (kill $(jobs -p)) and abort the pipeline.
    • Retry Loop Polish: Conditioned backoff sleep in process_shard to skip the 15s delay after the 3rd failed attempt.
  3. Validation & Configuration Updates:

    • Date Freshness SQL Validator: Added check_date_freshness_canonical_statvars to validation_config.json asserting MaxDate >= '2023' on active canonical StatVars (Count_Establishment_NAICSTransportationWarehousing, Count_Worker_NAICSRetailTrade).
    • Script-Based Import Alignment: Removed operational counters from manifest.json and shard_input_csv.sh. Removed legacy golden data dumps and omitted node_mcf since all StatVars are canonical/pre-resolved (0 new StatVar MCF nodes emitted).

Verification & Artifacts

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces parallel processing for downloading and processing Census County Business Patterns data, caches the StatVars MCF locally to optimize performance, and removes golden data validation checks. The review feedback highlights several critical improvement opportunities: streaming files directly in main.py to avoid potential Out-Of-Memory (OOM) errors during parallel execution, tracking and failing fast on sharding failures in Step 1 of shard_input_csv.sh, and ensuring proper error propagation inside the split_csv shell function.

Comment thread scripts/census_county_business_patterns/main.py Outdated
Comment thread scripts/census_county_business_patterns/shard_input_csv.sh Outdated
Comment thread scripts/census_county_business_patterns/shard_input_csv.sh
@kartik-s21
kartik-s21 force-pushed the fix-census-cbp-validation branch from 87e0fc4 to 34050f0 Compare September 11, 2026 11:07
@balit-raibot
balit-raibot self-requested a review September 17, 2026 03:23
balit-raibot and others added 3 commits September 17, 2026 03:24
…g and validation

- Fix [P1] empty shard pattern: add shopt -s nullglob, check for generated shards, and abort with exit 1 if no shards exist.
- Fix [P1] zero-observation check: verify observation data exists beyond CSV header via line count > 1.
- Fix [P1] operational counters: restore counters directory, pass --output_counters to stat_var_processor.py, clean up on failure, and add counters/* to manifest.json source_files.
- Fix [P2] runtime bottleneck: add safe bounded parallelism (default 8 workers) using bash wait -n job pool with fail-fast error propagation.
- Fix [P2] shell robustness: add set -euo pipefail, quote all variable expansions in split_csv, and verify split succeeds.
- Fix [P3] diagnostic logging: fix error message referencing SHARD_DIR instead of INPUT_DIR.
Comment on lines +109 to +116
if python3 "$STATVAR_PROCESSOR_SCRIPT" \
--input_data="$file" \
--existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf \
--pv_map="censuscountybusinesspatterns_pvmap.csv" \
--config_file="censuscountybusinesspatterns_metadata.csv" \
--output_path="$OUTPUT_FINAL_DIR/output_${prefix}" \
--counters_print_interval=-1
# --output_counters="$DEBUG_DIR/counters_${prefix}" \ # uncomment this line to debug the script like to get the details like memory utlization etc.
# Add any other required arguments for statvar_processpr.py here \
# Run in background

# Manage parallelism: pause if too many jobs are running
# We monitor the 'statvar_processpr.py' script's processes.
# sleep_while_active "$PARALLELISM" "$STATVAR_PROCESSOR_SCRIPT"
else
echo "WARNING: No shard files found matching '$INPUT_DIR/*_shard_*.csv' or '$file' is not a regular file."
echo "Please ensure your 'split_csv.sh' generates files in the '$INPUT_DIR' and follows the '*_shard_*.csv' naming convention."
--counters_print_interval=-1 \
--output_counters="$COUNTERS_DIR/counters_${prefix}.csv" && \

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Unverified Parallel Rewrite and Redundant 90 MB GCS Schema Downloads Across Workers

  • Finding: In commit 7bb12370, shard_input_csv.sh was updated from sequential execution to PARALLELISM=8 background workers, with each process_shard invocation passing --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf (89.88 MiB). Under util/file_util.py, stat_var_processor.py downloads gs:// inputs into a new /tmp/tmp* file on every invocation (and leaks the raw file descriptor from tempfile.mkstemp()). Across ~150 shards (and retries), this triggers >13 GB of redundant concurrent GCS downloads of the exact same 90 MiB file—the same source of transient GCS network/timeout failures that caused shards to fail originally. Furthermore, the Cloud Batch test run cited in the PR description (censuscountybusinesspatterns-samnotra-20260914-043508 / GCS version 2026_09_13T21_37_35_219006_07_00) ran on 2026-09-13 against commit 19a412b5 (sequential execution, before PARALLELISM=8 or "counters/*" were added).
  • Impact: High risk of GCS throttling, transient download timeouts, or /tmp disk pressure when 8 concurrent stat_var_processor.py workers repeatedly fetch stat_vars.mcf across ~150 shards, and the current head commit (7bb12370) has not been validated in a Cloud Batch test run.
  • Recommendation: Download gs://unresolved_mcf/scripts/statvar/stat_vars.mcf once to a local file (e.g., $SHARD_DIR/stat_vars.mcf) before Step 2 and pass --existing_statvar_mcf="$SHARD_DIR/stat_vars.mcf" to stat_var_processor.py. Then run a Cloud Batch test job on the updated commit and update the PR description with the new job link and GCS version URI (gs://datcom-import-test/...).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

addressed.

Comment on lines 14 to 19
"source_files": [
"gcs_output/source_files/*"
"gcs_output/source_files/*",
"counters/*"
],
"import_inputs": [
{

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Manifest Missing Mandatory node_mcf and golden_data/*.csv

  • Finding: manifest.json updates source_files to add "counters/*", but omits "golden_data/*.csv" even though validation_config.json configures GOLDENS_CHECK rules referencing ../../../../golden_data/golden_summary_report.csv and ../../../../golden_data/golden_observations.csv. Additionally, import_inputs only declares template_mcf and cleaned_csv (gcs_output/output/*.csv) and omits "node_mcf": "gcs_output/output/*.mcf", which is mandatory for stat_var_processor.py imports so generated output_*_stat_vars.mcf files are uploaded and resolved.
  • Impact: Golden CSV files are not archived to <version>/source_files/ in GCS as required by the golden validation guidelines, and any new StatVar MCF definitions emitted by stat_var_processor.py during future annual refreshes (e.g., new NAICS codes) will be dropped instead of ingested into the graph, causing Existence_MissingReference_variableMeasured validation failures.
  • Recommendation: Add "golden_data/*.csv" under source_files and add "node_mcf": "gcs_output/output/*.mcf" under import_inputs[0] in scripts/census_county_business_patterns/manifest.json.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

addressed.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Missing Date Freshness Validation and Unfiltered Production Dump in golden_summary_report.csv

  • Finding: While manifest.json configures "validation_config_file": "validation_config.json", inspecting the latest production summary report (2026_07_05T02_04_28_916907_07_00/input0/genmcf/summary_report.csv) shows 6,138 StatVars with non-uniform MaxDate values (2023: 5,988 StatVars; 2016: 141; 2022: 6; 2021: 3 due to historical NAICS revisions). However:
    1. validation_config.json contains no SQL_VALIDATOR date freshness rule verifying MaxDate (e.g., ensuring active/canonical StatVars or MAX(MaxDate) meet the expected 2023 release year within the 2-year CBP lag), leaving MaxDate unvalidated since golden_summary_report.csv excludes MaxDate.
    2. golden_data/golden_summary_report.csv checks in all 6,138 rows (444 KB) of the production summary_report.csv (including >6,000 auto-generated dc/... hash StatVars) instead of filtering against gs://unresolved_mcf/import_validation/nl_statvars.csv (which matches 71 StatVars), and golden_data/golden_observations.csv contains 10,391 rows (142 KB) instead of a concise 50–200 row slice.
  • Impact: Silent staleness if the latest year fails to download (especially since main.py:162-167 catches 404/BadZipFile on latest_year and continues without error), and excessive repo bloat / brittle golden checks from pinning NumPlaces across all 6,138 StatVars.
  • Recommendation:
    1. Add a scoped SQL_VALIDATOR rule in validation_config.json asserting freshness on MaxDate (e.g., verifying that canonical StatVars such as Count_Establishment_NAICSTransportationWarehousing / Count_Worker_NAICSRetailTrade have MaxDate >= '2023').
    2. Regenerate golden_data/golden_summary_report.csv filtered against gs://unresolved_mcf/import_validation/nl_statvars.csv (71 rows) using validator_goldens.py, and keep golden fixtures concise (50–200 rows).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

addressed.

Comment on lines +60 to +62
mkdir -p "$OUTPUT_FINAL_DIR"
mkdir -p "$SHARD_DIR"
echo "Directories created/ensured: $DEBUG_DIR, $OUTPUT_FINAL_DIR"
SHARD_ROWS=500000
mkdir -p "$COUNTERS_DIR"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Uncleaned Working Directories Across Repeated Executions and Delayed Failure Termination

  • Finding:
    1. Lines 60–62 run mkdir -p "$OUTPUT_FINAL_DIR" "$SHARD_DIR" "$COUNTERS_DIR" without clearing stale *_shard_*.csv, output_*, or counters_* files from prior executions. If shard_input_csv.sh is re-run in a persistent environment or after a partial run with fewer shards, shards=("$SHARD_DIR"/*_shard_*.csv) and "cleaned_csv": "gcs_output/output/*.csv" will pick up stale shard artifacts.
    2. In process_shard (lines 122–123), when attempt 3 of 3 fails, the script still logs "Retrying in 15s..." and sleeps 15 seconds before exiting the loop. Furthermore, when a background worker fails (job_failed=1 at line 148), the script breaks out of spawning new jobs but waits in while [ "$active_jobs" -gt 0 ]; do wait -n ... for all other active background workers to finish rather than terminating active background jobs immediately.
  • Impact: Potential ingestion of stale artifacts on repeated runs, and wasted Batch compute time on terminal shard failure.
  • Recommendation: Clean "$SHARD_DIR"/*_shard_*.csv, "$OUTPUT_FINAL_DIR"/output_*, and "$COUNTERS_DIR"/counters_* before Step 1; skip the sleep 15 backoff when attempt -eq 3; and terminate remaining background jobs (kill $(jobs -p) 2>/dev/null || true) when job_failed=1.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

addressed.

…nsusCountyBusinessPatterns

- Address code review comments:
  - Remove goldens data and node_mcf configuration for script-based import.
  - Remove operational counters generation from manifest and sharding script.
  - Update validation config with SQL date freshness check on canonical StatVars.
- Pipeline and Cloud Batch performance optimizations:
  - Parallelize Census zip downloads in main.py using ThreadPoolExecutor (8 threads).
  - Stream input files directly to DictReader in main.py, avoiding multi-GB memory buffers.
  - Parallelize text processing across CPU cores in main.py via ProcessPoolExecutor.
  - Cache schema MCF locally once in shard_input_csv.sh to eliminate redundant downloads.
  - Parallelize Step 1 sharding and inject headers via split --filter without sed -i disk rewrites.
  - Increase shard chunk size to 1M rows and scale PARALLELISM up to 48 workers on 64-core VMs.
  - Add fail-fast termination on worker job failure in shard_input_csv.sh.
@kartik-s21

Copy link
Copy Markdown
Contributor Author

Hi @rohitkumarbhagat,

Thank you for the thorough review. All feedback from Round 2 has been addressed in commit 7545935b, and the updated pipeline has been validated in Cloud Batch:

Summary of Fixes

  1. [P1] Local Schema Caching & Parallel Downloads (shard_input_csv.sh):

    • Download gs://unresolved_mcf/scripts/statvar/stat_vars.mcf once locally to $SHARD_DIR/stat_vars.mcf prior to Step 2, and pass --existing_statvar_mcf="$LOCAL_STATVAR_MCF" to workers. The local file is cleaned up after Step 2 completes.
    • This eliminated >13 GB of redundant GCS downloads across workers.
  2. [P2] Working Directory Hygiene & Fail-Fast Concurrency (shard_input_csv.sh):

    • Added pre-cleanup of $SHARD_DIR/*_shard_*.csv and $OUTPUT_FINAL_DIR/output_* before Step 1 to guarantee repeated runs start clean.
    • Attempt 3 skips the 15-second retry sleep.
    • Added immediate worker termination (kill $(jobs -p) 2>/dev/null || true) upon any worker failure in Step 1 or Step 2.
  3. [P2] Date Freshness SQL Validator (validation_config.json):

    • Added check_date_freshness_canonical_statvars asserting MaxDate >= '2023' for canonical StatVars (Count_Establishment_NAICSTransportationWarehousing, Count_Worker_NAICSRetailTrade). This passed cleanly in our test run.
  4. [P2] Manifest & Golden Data Configuration (manifest.json):

    • golden_data/*.csv: For script-based imports, we do not maintain static golden fixtures to avoid repo bloat and brittle checks. The legacy dump was removed, and data integrity is now enforced through the date freshness SQL check, check_empty_import, check_deleted_records_percent, and the automated version differ (which confirmed 0 diffs).
    • node_mcf: All StatVars in this import are canonical and defined via CensusCountyBusinessPatterns.tmcf and existing schema. stat_var_processor.py produces no new StatVar MCF nodes (output_*_stat_vars.mcf), confirmed by schema_diff_count: 0. Omitting node_mcf prevents glob resolution errors on non-existent files.
    • counters/*: Removed dead counter references from both manifest.json and shard_input_csv.sh.
  5. Additional Runtime Optimizations:

    • Streamed input CSVs directly to csv.DictReader and added ThreadPoolExecutor (8 workers) for zip downloads in main.py.
    • Parallelized Step 1 sharding with on-the-fly header injection via split --filter (avoiding sed -i rewrites).
    • Dynamically scaled worker pool up to 48 workers on 64-core VMs.

Cloud Batch Verification Results

  • Job Link: censuscountybusinesspatterns-samnotra-20260927-032649
  • Status: SUCCEEDED (Duration: 6h 03m, down from 18+ hours / timeouts)
  • GCS Version: gs://datcom-import-test/scripts/census_county_business_patterns/CensusCountyBusinessPatterns/2026_09_26T20_29_04_561640_07_00/
  • Validation:
    • check_date_freshness_canonical_statvars: PASSED
    • check_missing_refs_count: 0
    • check_lint_error_count: 0
    • check_deleted_records_percent: 0.0%
  • Differ Summary: obs_diff_count: 0 (73,457,237 observations matched baseline identically).

Please let me know if any further adjustments are needed!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants