Skip to content

Add a data pipeline: S3, Spark, PostgreSQL and Redshift - #5

Merged
vidit-16 merged 8 commits into
mainfrom
pipeline/postgres-spark-s3
Sep 24, 2026
Merged

vidit-16 merged 8 commits into
mainfrom
pipeline/postgres-spark-s3

Conversation

@vidit-16

@vidit-16 vidit-16 commented Sep 24, 2026 •

Copy link
Copy Markdown
Owner

Takes UCI Online Retail II from its source to two warehouses and a Power BI report, and checks every step against what the project already had.

What runs

  • Raw, curated and clean zones in a local folder or an S3 bucket.
  • jobs/conform_online_retail.py, a PySpark version of the Online Retail adapter. It imports nothing from the package, so the same file runs locally and on a cluster.
  • The existing cleaning and validation on Spark's output. Nothing is loaded while an error remains.
  • The SQL layer on PostgreSQL and Redshift Serverless, translated from the same files the SQLite tests run.
  • Report tables for BI (rate/mix split, country drivers, forecast accuracy), and a four-page Power BI report on them: docs/power_bi.md, docs/RootSignal.pbit.

Measured

  • Spark vs the pandas adapter on all 1,067,371 lines: all six tables identical, column by column.
  • 996,528 sales lines in PostgreSQL. Redshift loaded 996,531 from S3 with COPY, before stock codes were made case-insensitive; its seven queries matched SQLite.
  • Every staging view and all seven queries match SQLite row for row on PostgreSQL. aov may differ by a cent on exact half cents, since the engines round ties differently.
  • 427 tests. CI runs a PostgreSQL service and Java, so the warehouse and Spark tests run instead of skipping.
  • Mutation testing: 28 of 28 realistic pipeline bugs caught (22 before the gaps were tested), 93 of 108 operator mutants; the 15 left are wait times, output file counts and formatting.

Found and fixed along the way

  • 172 stock codes appear in two forms for one product (15056BL/15056bl, 47503J with a trailing space). Codes are now trimmed and upper-cased in both implementations.
  • Four queries sorted on a column that can tie, and SQLite and PostgreSQL put NULLs at opposite ends. Each query now sorts on its full grain, which a test enforces.
  • AVG of an integer column is an integer on Redshift. A test now requires every average to be of a real number.
  • Teardown reported removing a Glue job that never existed, because DeleteJob succeeds either way.
  • Spark wrote 200 files of about 100 KB for the sales fact; it now writes 4.

Not run yet
The Glue and EMR Serverless subcommands are written but have not run: the AWS account is not yet allowed to create Glue jobs, and its EMR Serverless quota had not taken effect. The Spark job has run locally and in the container, not on AWS. The README says so.

The staging views, mart and analytical queries are still written once. A
small translator rewrites the three SQLite constructs they use (week start,
year-month, REAL) and refuses anything it does not know, so a new query
fails at translation rather than on the server.

The warehouse loads the cleaned model with COPY inside one transaction, so a
load that breaks a foreign key leaves the previous warehouse untouched.

Tests load both engines from the same tables and require every staging view
and every query to match row for row. That found two things in the queries:
four sorted on a column that can tie, so the engines returned tied rows in
different orders, and SQLite and PostgreSQL place NULLs at opposite ends.
Each now sorts on a full key with NULL placement stated. Two edits make the
queries portable to Redshift as well: AVG of an integer column is cast
first, since Redshift returns an integer there, and the named WINDOW clause
is written out.

aov can differ by one cent on exact half-cent values: pandas rounds half to
even, SQLite rounds the binary float, PostgreSQL rounds the decimal half up.
Every other column matches exactly.
land writes the published workbook and a CSV rendering of it to the raw
zone. The lake is a local folder or an S3 bucket, laid out the same way.

conform is jobs/conform_online_retail.py, a PySpark version of the Online
Retail adapter. It imports nothing from the package, so the same file runs
under spark-submit, Glue and EMR Serverless. On all 1,067,371 published
lines it produces the same six tables as the pandas adapter, column by
column (scripts/compare_spark_to_pandas.py); tests hold the two to each
other on a fixture planting every defect the adapter handles.

load runs the existing cleaning and validation on Spark's output and
replaces the PostgreSQL warehouse only if no error remains: 3 errors before
cleaning, 2,590 rows quarantined, 996,531 sales lines loaded. It also
writes the cleaned model as Parquet typed to schema.sql, for Redshift.

publish writes the analysis results to the warehouse as report tables, so a
BI tool reads the project's findings rather than re-deriving them.
docs/power_bi.md describes the report built on them.

Spark runs in the project image, which now carries a Java runtime, because
Spark cannot write files on Windows without a Hadoop helper binary. CI
installs Java and runs a PostgreSQL service, so the Spark and warehouse
tests run there instead of skipping.
scripts/run_pipeline_aws.py runs the same stages with the lake in S3 and
Redshift Serverless as a second warehouse, one subcommand per step, so
nothing that bills is created without being named.

Run on 24 September 2026: the curated zone read from S3, the load wrote
the cleaned model to the S3 clean zone as Parquet, and Redshift loaded
996,531 sales lines from it with COPY. On the generated sample, all seven
analytical queries return exactly what SQLite returns.

Redshift is created at 8 RPU with a usage limit that deactivates it after
4 RPU-hours in a day, and teardown removes it. The bucket and IAM roles
are kept; they cost nothing.

The Glue and EMR Serverless subcommands are written but have not run: this
account is denied Glue, and its EMR Serverless vCPU quota was not yet in
effect.
Adds the pipeline's measured results, the path from S3 through Spark to
PostgreSQL and Redshift, and states that the Glue and EMR subcommands have
not run yet. Power BI is no longer out of scope; docs/power_bi.md covers
the report. Test count updated to 363.
tests/test_aws.py covers every AWS call the pipeline makes, with moto for
S3 and recording fakes for the rest: the roles reach only their bucket,
Glue and EMR get the lake paths and time limits, Redshift is created at one
size with a usage limit that deactivates it, statements run alone or as one
batch, failures and timeouts surface, and teardown stops before it deletes.

Writing them found that teardown reported removing a Glue job that never
existed, because DeleteJob succeeds either way. It now checks first.

tests/test_pipeline_stages.py pins the clean zone to schema.sql's types,
which Redshift's COPY requires, checks report tables flatten to plain
columns, and checks nothing loads while validation fails.

The image no longer copies data/lake in. The pipeline service mounts it,
and building it in added the raw CSV to every image.

docs/power_bi.md is now a step-by-step build of the report.
Four pages read the PostgreSQL warehouse: sales KPIs with a country slicer
and month-on-month change, each country's contribution to the latest
month's movement, the return rate split into rate and mix effects, and
forecast error against the naive baseline. Every figure comes from the fact
table or a report table the pipeline writes; the report derives nothing new.

docs/RootSignal.pbit is the report without data, and the screenshots are
rendered from its PDF export. Building it found that 172 product codes differ
only in case, which Power BI treats as duplicates; docs/power_bi.md records
it until the adapter handles it.
…g found

172 stock codes appear in two forms for one product, 15056BL and 15056bl,
or 47503J with a trailing space. 163 pairs carry the same description and
the rest the same product worded differently. Both the pandas adapter and
the Spark job now trim and upper-case the code, so a product's sales no
longer split across two rows, and Power BI no longer sees duplicate keys.

scripts/mutation_test_pipeline.py applies 28 realistic bugs one at a time.
Its first run caught 22; each miss had a cause, now covered:

- The Spark fixture lacked a credit note with a positive quantity, a
  product sold at both a zero and a real price, a non-breaking space, a
  zero quantity, and case variants of a code. It has them now.
- One Spark partition delivers rows in file order anyway, so taking the
  first value in any order went unnoticed. A test now feeds rows whose
  arrival order differs from their file order.
- No sample data ties in weekly movers, so a dropped tie-breaker passed. A
  test now requires every query to sort on its full output grain.
- Redshift averages an integer column to an integer; SQLite and PostgreSQL
  do not. A test now requires every AVG to average a cast or a REAL column.

The operator-level mutants also showed untested paths: the job's command
line entry point, a failed Spark run, reading the curated zone, nested S3
prefixes, exact CHECK removal, NULL numbers from PostgreSQL, the Redshift
query path, the timeout boundary and the spending defaults. Each has a test.

The backtest settings behind the README's 38.8% figure were written twice,
in the report tables and in scripts/run_external_dataset.py. Both now read
one constant.
…odes

The pipeline suite now catches all 28 realistic bugs in
scripts/mutation_test_pipeline.py, up from 22, and 93 of 108 operator
mutants; the 15 left are wait times, output file counts and formatting.

With stock codes trimmed and upper-cased, 172 duplicate products merge
(4,985 to 4,813) and PostgreSQL holds 996,528 sales lines. Redshift's
996,531 was loaded before the change and is reported as such. Every figure
in the Power BI report is unchanged, so its screenshots stand.
@vidit-16
vidit-16 merged commit 97a773b into main Sep 24, 2026
3 checks passed
@vidit-16
vidit-16 deleted the pipeline/postgres-spark-s3 branch September 24, 2026 21:27
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.

1 participant