Skip to content

Repository files navigation

VxIngest

VxIngest ingests meteorological data from various sources and prepares Couchbase-ready JSON documents for verification workflows used alongside the Model Application Toolsuite (MATS).

Getting Started

This repo currently ships the Python ingest application and related orchestration assets.

  • The ingest application reads GRIB2, NetCDF, or Couchbase source data and writes Couchbase-ready JSON documents to disk.
  • It also writes logs, Prometheus metrics, and tarballs in a transfer directory.
  • Those transfer tarballs are retained for downstream consumers, but the downstream import and metadata-update runtimes are not maintained in this branch.

If you want diagrams of the current data flow, see the Diagrams.

Usage

VxIngest is containerized for deployment. If you are developing the application, see docs/development-guide.md for environment setup, linting, formatting, and testing.

Using the container

Building images

VxIngest supports both AMD64 and ARM64 architectures.

Single-architecture local build:

docker build \
    --build-arg BUILDVER=dev \
    --build-arg COMMITBRANCH=$(git branch --show-current) \
    --build-arg COMMITSHA=$(git rev-parse HEAD) \
    -f ./docker/Dockerfile \
    -t vxingest/ingest:dev \
    .

Multi-architecture build with buildx:

docker buildx build \
    --platform linux/amd64,linux/arm64 \
    --build-arg BUILDVER=dev \
    --build-arg COMMITBRANCH=$(git branch --show-current) \
    --build-arg COMMITSHA=$(git rev-parse HEAD) \
    -f ./docker/Dockerfile \
    -t <registry>/vxingest/ingest:dev \
    --push \
    .

Build the development target used by the Compose test service:

docker build \
    --target dev \
    -f ./docker/Dockerfile \
    -t vxingest/ingest:dev-test \
    .

Running the ingest

Create a credentials file such as ${HOME}/credentials:

cb_host: "url.for.couchbase"
cb_user: "user"
cb_password: "password"
cb_bucket: "vxdata"
cb_scope: "_default"
cb_collection: "METAR"
cb_timeout_seconds: 7200

The optional cb_timeout_seconds sets Couchbase query timeouts. The cb_host value must include a protocol such as couchbase:// or couchbases://.

Run the ingest through Docker Compose. The Compose service already supplies the standard output, log, metrics, and transfer directories; you only need to provide the job identifier:

data=/data-ingest/data \
public=/public \
docker compose run ingest \
    -j JOB-TEST:V01:METAR:NETCDF:OBS

The ingest writes JSON output, logs, metrics, and transfer tarballs into the mounted data directory.

Import lock and readiness coordination

When VxIngest runs a builder that reads model or observation documents written by VxImporter, the builder checks the shared import lock before reading the dataset. The lock is stored in the target bucket's default scope and COMMON collection under this document key:

MD:import_lock:COMMON:V01

VxImporter owns this document. During an import it writes status="running", refreshes the Unix-seconds updated field every minute, and records an identifying job_id containing its host name and process ID. It changes the status to idle during normal deferred cleanup. A typical document is:

{
    "id": "MD:import_lock:COMMON:V01",
    "status": "running",
    "updated": 1760000000,
    "job_id": "vximporter:hostname:12345"
}

When two VxImporter jobs start close together, the later job waits for the fresh running lock instead of starting concurrently. It polls every 10 seconds for up to 30 minutes. A heartbeat older than 30 minutes is treated as stale and may be reclaimed by the waiting importer with a CAS-protected replacement.

The VxIngest reader behavior is:

  • Poll the lock every 10 seconds while status is running.
  • Proceed immediately when the document is absent or its status is not running.
  • Treat a running lock as stale when updated is more than 30 minutes old, then proceed and log a warning.
  • Proceed after waiting 30 minutes even if the lock remains fresh.
  • Proceed when the lock cannot be read; lock checking is fail-open so a Couchbase read error does not block a builder indefinitely.

This wait prevents normal CTC and partial-sums reads from observing a dataset while VxImporter is still writing it. It does not provide a transaction or rollback: a stale-lock decision means the builder may read a partially imported dataset. Operators should inspect the lock's job_id and updated fields and the importer logs before treating a stale lock as harmless.

The lock wait is separate from Couchbase bucket readiness. The ingest application's Couchbase clients use their configured SDK timeouts when opening connections; the VxImporter process waits for its bucket with a default 60-second BUCKET_READY_TIMEOUT_SECONDS value. Increasing that VxImporter environment variable helps slow or remote clusters become ready, but it does not change the VxIngest reader poll, stale-lock, or maximum-wait values.

Testing mode

To run ingest in testing mode, set the TESTING environment variable (any value) when running the container. When set, the ingest will process both status='active' and status='test' job documents. When not set, only status='active' documents are processed. This allows test documents to be safely developed and tested without risk of automatic runners (like cron) inadvertently executing them:

data=/data-ingest/data \
public=/public \
docker compose run -e TESTING=1 ingest \
    -j JOB-TEST:V01:METAR:NETCDF:OBS

Running tests in the container

The Compose test service builds the Dockerfile's dev target and runs the repository test suite inside that container:

data=/home/path/to/test-data docker compose run test

Running jobs with the job wrapper script

For production or automation workflows, use scripts/VXingest_utilities/run_job.sh to submit and process ingest jobs with Docker and automatically import the resulting documents into Couchbase.

The script:

  • Orchestrates ingest and import of a single job
  • Manages temporary working directories and logs
  • Automatically extracts tar.gz archives from ingest output and imports any JSON documents found within
  • Handles both ingest output and Couchbase document import
  • Requires the CREDENTIALS_FILE environment variable to point to a credentials YAML file (see Running the ingest)

Basic usage:

export CREDENTIALS_FILE="${HOME}/credentials"
./scripts/VXingest_utilities/run_job.sh JOB-TEST:V01:METAR:NETCDF:OBS

Optional environment variables:

  • WORKING_ROOT_DIR — Root directory for temporary files and logs. Default: /data-ingest/data/working
  • PUBLIC_DIR — Host public directory mounted into ingest as /public. Default: /public
  • DATA_SOURCE — Host directory containing raw input files. If set, this directory is mounted read-only into the container. Optional.
  • CONTAINER_DATA_PATH — Container path where DATA_SOURCE is mounted. Default: same as DATA_SOURCE. Only used if DATA_SOURCE is set.
  • VXINGEST_IMAGE — Docker image for the ingest step. Default: ghcr.io/noaa-gsl/vxingest/ingest:latest
  • LOG_LEVEL — Log level for both ingest and import steps. Use one of DEBUG, INFO, WARNING, ERROR, or CRITICAL. Default: INFO
  • VXIMPORTER_IMAGE — Docker image for the import step. Default: ghcr.io/noaa-gsl/vximporter:latest
  • VXIMPORTER_WORKERS — Number of import workers. Default: 16
  • VXIMPORTER_BATCH_SIZE — Batch size for imports. Default: 1000
  • VX_METADATA_UPDATER_SETTINGS — Optional host path to the VxMetadataUpdater settings.json file. When set, it is mounted into the updater container at /app/settings.json and passed via -s /app/settings.json. When unset, the container's default settings are used.
  • VX_METADATA_UPDATER_IMAGE — Docker image for the metadata update step in ingest_model.sh. Default: ghcr.io/noaa-gsl/vxmetadataupdater:latest
  • VX_METADATA_UPDATER_DOCKER_USER — Optional user/group override for the metadata updater Docker run. Default: unset, so the image's default user is used.

The script creates logs in ${WORKING_ROOT_DIR}/logs/ with naming pattern docker-{ingest|import}-{job_id}-{timestamp}.out.

Example with data source:

export CREDENTIALS_FILE="${HOME}/credentials"
export DATA_SOURCE="/opt/data/netcdf_to_cb"
./scripts/VXingest_utilities/run_job.sh JS:METAR:OBS:NETCDF-TEST:schedule:job:V01

Example with debug logging:

export CREDENTIALS_FILE="${HOME}/credentials"
export LOG_LEVEL=DEBUG
./scripts/VXingest_utilities/run_job.sh JS:METAR:OBS:NETCDF-TEST:schedule:job:V01

Using Docker Compose directly

The wrapper script uses direct docker run calls so it is self-contained for automation. Docker Compose remains supported for development, testing, and interactive debugging through compose.yaml. Use Compose when you want the repository-defined shell, test, or ingest services rather than the wrapper's ingest-plus-import workflow.

Example direct Compose ingest run:

data=/data-ingest/data \
public=/public \
LOG_LEVEL=DEBUG \
docker compose run ingest \
    -j JOB-TEST:V01:METAR:NETCDF:OBS

LOG_LEVEL controls application logging for the main process and worker processes. If it is unset, VxIngest logs at INFO. Valid values are DEBUG, INFO, WARNING, ERROR, and CRITICAL. Invalid values stop startup with an error so misconfigured automation does not silently run at the wrong verbosity.

Debugging in the container

If you want an interactive shell in the ingest image for debugging:

data=/data-ingest/data \
public=/public \
docker compose run shell

Debugging with a locally built image

To debug VxIngest with a locally built Docker image and DEBUG logging:

  1. Build the image locally:
docker build \
    --build-arg BUILDVER=dev \
    --build-arg COMMITBRANCH=$(git branch --show-current) \
    --build-arg COMMITSHA=$(git rev-parse HEAD) \
    -f ./docker/Dockerfile \
    -t vxingest/ingest:latest \
    .
  1. Set the log level and image environment variables:
export LOG_LEVEL=DEBUG
export VXINGEST_IMAGE=vxingest/ingest:latest
export CREDENTIALS_FILE="${HOME}/credentials"
  1. Run the job wrapper script:
./scripts/VXingest_utilities/run_job.sh JS:METAR:OBS:NETCDF-TEST:schedule:job:V01

The DEBUG logs will be written to ${WORKING_ROOT_DIR}/logs/ (default: /data-ingest/data/working/logs/). Check the log files for detailed output including the get_file_list debug entries showing which files were discovered and why.

Note on reprocessing files: If you need to reprocess files that have already been ingested, you may need to DELETE type "DF" documents from the database that record the file's processing status. The ingest code will NOT process any files that have been recorded in type "DF" documents indicating they have already been processed. Use a query like this to remove the processing record, the extra fields are there to cause the query to make use of proper indexing:

DELETE
FROM vxdata._default.METAR
WHERE subset='METAR'
    AND type='DF'
    AND fileType='netcdf'
    AND originType='madis'
    AND url = "/opt/data/netcdf_to_cb/input_files/20250911_1500"

Alternatively, you can touch the input data file i.e. ...

touch /opt/data/netcdf_to_cb/input_files/20250911_1500

... which will renew the mtime of the input file making it more recent than the previously ingested data. This will allow reprocessing as well.

Tailing the log output from a running contianer

The log files are identified in the output of of the run_job.sh. You can tail these in real time with ...

tail -f log_file

If you want to tail the log output from the latest running container use...

docker logs -f "$(docker ps -ql)"

which will tail the most recent container. Alternatively you can examine the running containers with ...

docker ps

and choose a running container (in case there are more than one running container) and then use ...

docker logs -f container_id

Diagrams

Data flow for model and observation ingest (GRIB2 and NetCDF):

---
title: Model and Obs Ingest
---
flowchart LR
    data --> |1. Reads new data| ingest
    ingest --> |2. Writes data out as JSON files| disk
    disk --> |3. Hands tarballs to downstream tooling| downstream
    downstream --> |4. Inserts files| cb

    subgraph Application Layer
        ingest(Ingest)
        downstream(External downstream import)
    end
    subgraph Data Layer
        data[[Model and Obs Data]]
        disk[[Files on Disk]]
        cb[(Couchbase)]
    end
Loading

Data flow for aggregate statistics ingest (CTC and Partial Sums):

---
title: CTC and Partial Sums Ingest
---
flowchart LR
    ingest --> |1. Gets data from Couchbase| cb
    ingest --> |2. Writes data out as JSON files| disk
    disk --> |3. Hands tarballs to downstream tooling| downstream
    downstream --> |4. Inserts files| cb

    subgraph Application Layer
        ingest(Ingest)
        downstream(External downstream import)
    end
    subgraph Data Layer
        disk[[Files on Disk]]
        cb[(Couchbase)]
    end
Loading

Disclaimer

This repository is a scientific product and is not official communication of the National Oceanic and Atmospheric Administration, or the United States Department of Commerce. All NOAA GitHub project code is provided on an "as is" basis and the user assumes responsibility for its use. Any claims against the Department of Commerce or Department of Commerce bureaus stemming from the use of this GitHub project will be governed by all applicable Federal law. Any reference to specific commercial products, processes, or services by service mark, trademark, manufacturer, or otherwise, does not constitute or imply their endorsement, recommendation or favoring by the Department of Commerce. The Department of Commerce seal and logo, or the seal and logo of a DOC bureau, shall not be used in any manner to imply endorsement of any commercial product or activity by DOC or the United States Government.

About

Extract, transform, and load model and observation data for model verification

Resources

Stars

2 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages