Metadata-Version: 2.5
Name: sws-spark-dissemination-helper
Version: 0.1.18
Summary: A Python helper package providing streamlined Spark functions for efficient data dissemination processes
Project-URL: Repository, https://github.com/un-fao/fao-sws-it-python-spark-dissemination-helper
Author-email: Daniele Mansillo <danielemansillo@gmail.com>, Luca Simi <luca.simi@fao.org>
License: MIT License
        
        Copyright (c) 2024 Daniele Mansillo
        
        Permission is hereby granted, free of charge, to any person obtaining a copy
        of this software and associated documentation files (the "Software"), to deal
        in the Software without restriction, including without limitation the rights
        to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
        copies of the Software, and to permit persons to whom the Software is
        furnished to do so, subject to the following conditions:
        
        The above copyright notice and this permission notice shall be included in all
        copies or substantial portions of the Software.
        
        THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
        IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
        FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
        AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
        LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
        OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
        SOFTWARE.
License-File: LICENSE
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Requires-Python: >=3.9
Requires-Dist: boto3>=1.40.0
Requires-Dist: botocore>=1.40.0
Requires-Dist: httpx==0.28.1
Requires-Dist: pyspark==3.5.6
Requires-Dist: python-dotenv==0.19.2
Requires-Dist: sws-api-client<3,>=2.18.0; python_version == '3.11'
Requires-Dist: sws-api-client==2.7.3; python_version == '3.9'
Provides-Extra: dev
Requires-Dist: black==24.10.0; extra == 'dev'
Requires-Dist: build>=1.0.0; extra == 'dev'
Requires-Dist: pyarrow>=4.0.0; extra == 'dev'
Requires-Dist: pytest-cov; extra == 'dev'
Requires-Dist: pytest>=8.0.0; extra == 'dev'
Requires-Dist: twine>=4.0.0; extra == 'dev'
Provides-Extra: docs
Requires-Dist: furo; extra == 'docs'
Requires-Dist: sphinx>=7.0; extra == 'docs'
Description-Content-Type: text/markdown

# sws-spark-dissemination-helper

Python/PySpark helper library for the FAO SWS (Statistical Working System) data dissemination pipelines. It provides a set of classes that move a dataset from the SWS Postgres database, through an Iceberg-backed Bronze → Silver → Gold medallion pipeline, to the CSV/Iceberg artifacts and SWS "dissemination tags" consumed downstream (including SDMX and FAOSTAT-shaped outputs).

## Pipeline overview

| Layer | Class | Role |
| --- | --- | --- |
| Source | [`SWSPostgresSparkReader`](src/sws_spark_dissemination_helper/SWSPostgresSparkReader.py) | Reads SWS Postgres tables (observations, codelists, datatables) into Spark and stages them as Iceberg tables. Every stage reads the specific datatables it needs on demand through this class, rather than a single upfront import. |
| Bronze | [`SWSBronzeIcebergSparkHelper`](src/sws_spark_dissemination_helper/SWSBronzeIcebergSparkHelper.py) | Denormalizes the raw observation/coordinate/metadata tables into one wide table (codes instead of ids, metadata and unit-of-measure attached). Can also be filtered to a dimension-value selection provided as input and republished as a separate "disseminated tag" table. |
| Silver | [`SWSSilverIcebergSparkHelper`](src/sws_spark_dissemination_helper/SWSSilverIcebergSparkHelper.py) | Applies dissemination business rules on top of Bronze data: code correction, time-validity checks, dissemination flags/exceptions, duplicate detection. |
| Gold | [`SWSGoldIcebergSparkHelper`](src/sws_spark_dissemination_helper/SWSGoldIcebergSparkHelper.py) | Produces the published output shapes from Silver/Bronze data: SWS variants, a pre-SDMX/FMR-layout variant (SDMX doesn't accept dots in coded values and requires the value column named `OBS_VALUE`), an SDMX-shaped variant, and FAOSTAT variants — the SDMX and FAOSTAT variants are both produced by mapping the pre-SDMX output through a remote SDMX structure-map chain run as SWS plugin tasks. Also applies display-decimal rounding. |
| Standalone | [`SWSEasyIcebergSparkHelper`](src/sws_spark_dissemination_helper/SWSEasyIcebergSparkHelper.py) | Exports a huge dataset to Iceberg as a single denormalized table, for datasets large enough that querying them directly in SWS is impractical or impossible. |

Each helper writes its output to Iceberg, creates a version tag, optionally caches a CSV copy on S3, and can register the result in the SWS dissemination-tag metadata via the `sws_api_client` `Tags` API. See each class's docstring for the exact tables/columns it reads and writes.

## Installation

```sh
pip install sws-spark-dissemination-helper
```

Requires Python 3.9 or 3.11 (see `pyproject.toml`).

Construct `SWSPostgresSparkReader` once per job, then pass it into whichever helper class(es) you need for the pipeline stage(s) you're running — see the [architecture page](docs/source/architecture.rst) for how the stages relate and the [API reference](docs/source/api/index.rst) (or each class's docstring) for the exact constructor arguments.

## Configuration

The package expects to run on Amazon EMR, with AWS access already configured (an EMR instance profile):

- `EMR_BUCKET`: S3 bucket used as the Spark/Iceberg warehouse location (see `utils.get_spark`).
- `DB_SECRET`: name of the AWS Secrets Manager secret holding the Postgres JDBC connection properties (`host`, `port`, `database`, `user`, `password`), used when `SWSPostgresSparkReader` is instantiated without explicit `jdbc_conn_properties`.
- AWS credentials are picked up from the default `boto3` credential chain and used both for S3/Glue access and to build the temporary Spark `s3a` credentials.

## Documentation

The full documentation (architecture, configuration, and the generated API reference) is built with Sphinx from this README and the package's docstrings. Build it locally with:

```sh
pip install -e ".[docs]"
make docs
```

Then open `docs/build/html/index.html`. See [docs/source/](docs/source/) for the source.

## How to use (development)

Make sure to have `make` installed on your system and run `make help` to see the available commands.

## Upload a new version

A new version is published to PyPI automatically by the [publish workflow](./.github/workflows/publish.yml) whenever a GitHub release is published. The workflow verifies that the release tag matches `v<project.version>`, so the two must always be kept in sync.

1. Update `project.version` in the [pyproject.toml](./pyproject.toml) file (e.g. `0.1.9.dev0`) and commit the change.
2. Create and push a tag matching the new version (note the leading `v`):

   ```sh
   git tag v0.1.9.dev0
   git push origin --tags
   ```

3. Go to the [tags page](https://github.com/un-fao/fao-sws-it-python-spark-dissemination-helper/tags), select the tag and create a release from it. Use the **Latest** release label for official releases or **Pre-release** for development versions, then fill in the release notes.
4. Wait for the GitHub action to run. If it succeeds the new version will be available on PyPI; otherwise check the workflow logs for the failure reason (a common one is the tag not matching `project.version`) and fix it.
