An end-to-end data architecture for a freight forwarder / 3PL in global ocean container logistics. Real maritime data (AIS vessel tracking, port/vessel/carrier reference, trade-flow priors) plus seeded synthetic events flow through a GCP pipeline into two analytical stores: a BigQuery star-schema warehouse for OLAP questions and an ArangoDB property graph for network questions. Each of the four use cases runs on whichever store fits the workload.
A 3-person group project: take a real domain, turn it into a data model, build a cloud ETL pipeline, orchestrate it with Airflow, and demo it end to end.
| # | Question | Store | Why that store |
|---|---|---|---|
| UC1 | ETA reliability — actual vs. proforma schedule delta | BigQuery | Aggregation over a partitioned fact; columnar scan |
| UC2 | Port congestion / dwell over time | BigQuery | Time-series rollup on a date-partitioned fact |
| UC3 | Chokepoint exposure (Suez / Panama / Malacca …) | ArangoDB | Centrality + reachability over the route network |
| UC4 | Rerouting / alternative-path analysis | ArangoDB | Weighted shortest-path traversal |
OLAP questions run on the warehouse and network questions run on the graph. The same business keys join a graph result back to a warehouse fact.
flowchart LR
subgraph SRC[Sources]
AIS[AIS vessel tracking - real]
REF[Port / vessel / carrier reference - real]
PRIOR[LSCI / LPI / Comtrade priors - real]
SYN[Synthetic schedules, bookings, events - seeded]
end
SRC --> BRONZE[(GCS Bronze<br/>raw, immutable<br/>Parquet + JSONL)]
BRONZE --> DAG[Airflow 3.0 DAG<br/>ofa_warehouse]
DAG --> SILVER[(GCS Silver<br/>conformed staging<br/>one source of truth)]
SILVER --> BQ[(BigQuery<br/>ofa_star — star schema)]
SILVER --> AG[(ArangoDB<br/>ocean_network graph)]
BQ --> UC12[UC1 ETA reliability<br/>UC2 congestion]
AG --> UC34[UC3 chokepoint exposure<br/>UC4 rerouting]
Both stores are built from one conformed Silver layer. The business keys (UN/LOCODE for
ports, IMO for vessels, SCAC for carriers) serve as both the BigQuery dimension natural
keys and the ArangoDB _keys, so a graph traversal result lines up 1:1 with a warehouse
row.
flowchart TB
B[Bronze — raw landing<br/>ais/dt= · reference/ · priors/ · synthetic/] -->|stage_conform| S[Silver — conformed<br/>dim_*.parquet · fact_*/dt=]
S -->|BQ load leg| W[BigQuery ofa_star]
S -->|graph load leg| G[ArangoDB ocean_network]
This project runs plain Apache Airflow 3.0, not Cloud Composer. The DAG uses only
standard operators so it stays Composer-portable: it runs locally with
airflow dags test today and lifts onto Cloud Composer 3 unchanged if the project moves to
managed Airflow. One conform step produces Silver. The BigQuery and ArangoDB load legs run
in parallel off that Silver, and verify fans in only after both stores finish.
flowchart TD
CONFORM[stage_conform<br/>build Silver]
subgraph BQLEG[BigQuery load leg]
direction TB
STGV[load_staging_dim_vessel] --> MRGV[merge_dim_vessel<br/>SCD2 upsert]
STGC[load_staging_dim_carrier] --> MRGC[merge_dim_carrier<br/>SCD2 upsert]
OWP[overwrite_dim_port<br/>SCD1]
OWL[overwrite_dim_lane<br/>SCD1]
OWO[overwrite_operated_by<br/>bridge]
OWVL[overwrite_fact_voyage_leg<br/>partition + cluster]
OWPC[overwrite_fact_port_call<br/>partition + cluster]
end
ARANGO[load_arango<br/>idempotent UPSERT]
VERIFY[verify<br/>row counts + cross-store reconcile<br/>+ live UC3/UC4 anti-degeneracy]
CONFORM --> STGV & STGC & OWP & OWL & OWO & OWVL & OWPC & ARANGO
OWP & OWL & OWO & OWVL & OWPC & MRGV & MRGC & ARANGO --> VERIFY
verify fans in on both legs because the cross-store checks can only run once BigQuery and
ArangoDB are both populated. A half-loaded store can't pass or fail the gate by accident.
Two date-partitioned, clustered facts surrounded by flat conformed dimensions. We use a
star rather than a snowflake: on a columnar engine storage is cheap and joins are the cost,
so normalizing dimensions to save space is the wrong trade (full rationale in
docs/deck/m2-star-vs-snowflake.md).
erDiagram
dim_vessel ||--o{ fact_voyage_leg : vessel_imo
dim_port ||--o{ fact_voyage_leg : origin_unlocode
dim_port ||--o{ fact_voyage_leg : dest_unlocode
dim_lane ||--o{ fact_voyage_leg : lane_key
dim_vessel ||--o{ fact_port_call : vessel_imo
dim_port ||--o{ fact_port_call : unlocode
dim_vessel ||--o{ operated_by : vessel_imo
dim_carrier ||--o{ operated_by : carrier_scac
fact_voyage_leg {
string vessel_imo FK
string origin_unlocode FK
string dest_unlocode FK
float transit_hours
float distance_nm
float schedule_delta
date dt "partition"
}
fact_port_call {
string vessel_imo FK
string unlocode FK
timestamp arrival_ts
timestamp departure_ts
date dt "partition"
}
dim_vessel {
string imo PK "SCD2"
string vessel_name
date effective_from
date effective_to
bool is_current
}
dim_carrier {
string scac PK "SCD2"
string carrier_name
bool is_current
}
dim_port {
string unlocode PK "SCD1"
float lat
float lon
}
dim_lane {
string lane_key PK "SCD1"
string origin_unlocode
string dest_unlocode
}
operated_by {
string vessel_imo FK
string carrier_scac FK
}
DDL is in sql/ddl_star.sql; the SCD2 dimension upserts are
sql/merge_dim_vessel.sql and
sql/merge_dim_carrier.sql. Facts partition on dt and cluster
on the foreign keys queried most.
Five vertex collections, four edge collections, one named graph. The model is fixed in
docs/deck/m2-arango-graph.md.
flowchart LR
V[vessels<br/>_key = IMO]
C[carriers<br/>_key = SCAC]
P[ports<br/>_key = UN/LOCODE]
L[lanes<br/>_key = origin__dest]
K[chokepoints<br/>_key = code]
C -->|operates| V
V -->|calls_at| P
P -->|route| P
L -->|transits_chokepoint| K
The same key identifies a thing in both stores, so the two halves share one identity space.
flowchart LR
subgraph BQ[BigQuery ofa_star]
BV[dim_vessel.imo]
BP[dim_port.unlocode]
BC[dim_carrier.scac]
BL[dim_lane.lane_key]
end
subgraph AG[ArangoDB ocean_network]
AV[vessels._key]
AP[ports._key]
AC[carriers._key]
AL[lanes._key]
end
BV --- AV
BP --- AP
BC --- AC
BL --- AL
| Source | Real / synthetic | Format | What it gives us |
|---|---|---|---|
| AIS vessel tracking (MarineCadastre) | Real | GeoParquet | Vessel positions → port calls and voyage legs |
| Port / vessel / carrier reference (UN/LOCODE, World Port Index) | Real | CSV → Parquet | Dimension attributes, geocoding |
| Trade-flow priors (UNCTAD LSCI, World Bank LPI, UN Comtrade) | Real | JSON/CSV | Condition the synthetic lane network |
| Schedules, bookings, container events | Synthetic (seeded) | JSONL | Proforma timing, booking/event volume |
| Lane network + chokepoint assignments | Synthetic (curated) | constants | The route graph and chokepoint transits |
Real vs. synthetic. AIS and reference data are real but limited to a bounded slice
(Q1 2024, four US ports: Houston, LA, NYC, Savannah) to keep scan cost and graph size in
check. The booking/event/schedule layer and the foreign-port lane network are synthetic, since
no single public feed carries a forwarder's bookings end to end. We generate them
deterministically so any teammate reproduces byte-identical output. Generators are seeded
with numpy.random.default_rng and a pinned Faker version, and the contract is recorded in
synthetic.sha256. samples/ holds small proof-of-pull extracts from
each real source so the sourcing is visible without committing bulk data.
Everything runs through make. Targets shell to python -m <module> or to bq / airflow.
# 1. Install (Python 3.11/3.12). Airflow installs via its constraints file.
pip install -e .[dev]
# 2. Credentials. Copy the template and fill in the real (gitignored) .env.
cp .env.template .env # GCP via ADC; ARANGO_* point at the managed cluster
# 3. Bronze — pull real sources, generate synthetic, land in GCS, verify.
make bronze # = pull-reference pull-priors pull-ais generate load-bronze verify
# 4. Silver — conform + derive the staging layer, verify.
make silver # = conform derive verify
# 5. Warehouse — create the star schema, load it, verify.
make warehouse # = ddl load-bq verify (runs the DAG via `airflow dags test`)
# 6. Graph — load the ArangoDB named graph.
make load-arango
# 7. Freeze the four UC answers and run the demo notebook end-to-end.
make freeze
make demoUseful focused targets:
| Target | What it does |
|---|---|
make verify |
Full verification ladder (Bronze → Silver → BQ → cross-store → UC) |
make verify-uc |
Live UC3/UC4 anti-degeneracy check against the ArangoDB cluster |
make verify-cluster |
Connectivity + load smoke test against the cluster |
make secret-gate |
No-committed-secrets scan (same check CI runs) |
The design is done and a working vertical slice runs end to end.
- Pipeline — Bronze ingestion, Silver conform/derive, BigQuery star load, and ArangoDB
graph load all run; the
ofa_warehouseDAG wires them with a fan-inverify. - Modeling — ER, dimensional, and graph design are complete (see
docs/). - Demo — all four use cases answer;
docs/demo.ipynbruns from a clean clone against frozen golden snapshots, so the demo holds with no live credentials.
| Layer | Choice |
|---|---|
| Raw + staging | GCS, Hive-style date partitioning; Parquet for tabular, JSONL for events |
| Orchestration | Plain Apache Airflow 3.0, Composer-portable |
| OLAP store | BigQuery — native tables, star schema, partitioned + clustered |
| Graph store | Managed ArangoDB 3.12 cluster |
| Graph analytics | AQL traversal / SHORTEST_PATH + Graph Analytics Engine (GAE) server-side; NetworkX is the fallback |
| Synthetic data | Faker (pinned) + NumPy default_rng + pandas/pyarrow |
| Directory | Purpose |
|---|---|
ingest/ |
Real-source pulls — AIS, reference, priors |
data_gen/ |
Seeded synthetic generators — schedules, bookings, container events, lane network |
silver/ |
Conform + derive: identity resolution, geofencing, dims, facts |
sql/ |
BigQuery DDL, SCD2 merges, UC1/UC2 query SQL |
aql/ |
ArangoDB queries for UC3/UC4 (with .explain plans) |
analytics/ |
UC3/UC4 runners, criticality, snapshot builders |
lib/ |
Shared clients + loaders — GCS, ArangoDB, graph loader/queries |
dags/ |
The ofa_warehouse Airflow DAG |
scripts/ |
Bronze landing, verify checks, UC freeze, demo build |
tests/ |
pytest suite (conform, derive, geofence, graph load, cross-store, DAG) |
docs/ |
Design notes (deck/), demo notebook, run-from-clone guide |
reference/ |
Small committed reference tables (e.g. chokepoints) |
samples/ |
Proof-of-pull extracts from each real source |
web/ |
Next.js demo app (live BigQuery with golden fallback) |
Design notes live as Markdown in docs/deck/; the Mermaid renders on GitHub.
| Doc | Covers |
|---|---|
| er-logical.md | Logical ER + fact/dimension classification |
| bq-star.md | Star: grains, conformed dims, SCD, partition/cluster |
| star-vs-snowflake.md | Star-over-snowflake rationale |
| arango-graph.md | Property-graph model, named graph ocean_network |
| conformed-keys.md | Conformed-key bridge + hybrid justification |
| gap-analysis.md | Data-needs vs. available-sources reconciliation |
Domain and use-case notes (m1-team-domain.md, m1-use-cases.md, m1-source-inventory.md,
m1-real-vs-synthetic.md, m1-billing-guard.md) are in the same folder.