Skip to content

Commit b6e27f3

Browse files
committed
feat: implement direct GCS streaming by URIs via BigQuery metadata; update docs and unit tests accordingly
1 parent 54fdefa commit b6e27f3

17 files changed

Lines changed: 668 additions & 164 deletions
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
view,timestamp,logger,memory,unit
2+
DEFAULT,2026-04-20T14:59:18.784422Z,METRIC_MEM:,429,MB
3+
DEFAULT,2026-04-20T14:59:19.784315Z,METRIC_MEM:,467.35,MB
4+
DEFAULT,2026-04-20T14:59:20.784602Z,METRIC_MEM:,525.09,MB
5+
DEFAULT,2026-04-20T14:59:21.784888Z,METRIC_MEM:,530.13,MB
6+
DEFAULT,2026-04-20T14:59:22.785227Z,METRIC_MEM:,533.62,MB
7+
DEFAULT,2026-04-20T14:59:23.785545Z,METRIC_MEM:,534.02,MB
8+
DEFAULT,2026-04-20T14:59:24.785893Z,METRIC_MEM:,534.46,MB
9+
DEFAULT,2026-04-20T14:59:25.786286Z,METRIC_MEM:,533.57,MB
10+
DEFAULT,2026-04-20T14:59:26.786608Z,METRIC_MEM:,533.07,MB
11+
DEFAULT,2026-04-20T14:59:27.786989Z,METRIC_MEM:,533.24,MB
12+
DEFAULT,2026-04-20T14:59:28.787246Z,METRIC_MEM:,534.23,MB
13+
DEFAULT,2026-04-20T14:59:29.787577Z,METRIC_MEM:,534.72,MB
14+
DEFAULT,2026-04-20T14:59:30.787938Z,METRIC_MEM:,534.72,MB
15+
DEFAULT,2026-04-20T14:59:31.788315Z,METRIC_MEM:,534.96,MB
16+
DEFAULT,2026-04-20T14:59:32.788644Z,METRIC_MEM:,536.2,MB
17+
DEFAULT,2026-04-20T14:59:33.788993Z,METRIC_MEM:,538.16,MB
18+
DEFAULT,2026-04-20T14:59:34.789341Z,METRIC_MEM:,538.66,MB
19+
DEFAULT,2026-04-20T14:59:35.789670Z,METRIC_MEM:,538.66,MB
20+
DEFAULT,2026-04-20T14:59:36.789550Z,METRIC_MEM:,539.64,MB
21+
DEFAULT,2026-04-20T14:59:37.789398Z,METRIC_MEM:,539.64,MB
22+
DEFAULT,2026-04-20T14:59:38.789317Z,METRIC_MEM:,541.61,MB
23+
DEFAULT,2026-04-20T14:59:39.789179Z,METRIC_MEM:,541.61,MB
24+
DEFAULT,2026-04-20T14:59:40.789326Z,METRIC_MEM:,546.59,MB
25+
DEFAULT,2026-04-20T14:59:41.789650Z,METRIC_MEM:,546.59,MB
26+
DEFAULT,2026-04-20T14:59:42.790023Z,METRIC_MEM:,546.59,MB
27+
DEFAULT,2026-04-20T14:59:43.790569Z,METRIC_MEM:,720.57,MB
28+
DEFAULT,2026-04-20T14:59:44.790956Z,METRIC_MEM:,950.8,MB
29+
DEFAULT,2026-04-20T14:59:45.791405Z,METRIC_MEM:,1116.99,MB
30+
DEFAULT,2026-04-20T14:59:46.791919Z,METRIC_MEM:,1334.91,MB
31+
DEFAULT,2026-04-20T14:59:47.792322Z,METRIC_MEM:,1531.51,MB
32+
DEFAULT,2026-04-20T14:59:48.792605Z,METRIC_MEM:,1672.36,MB
33+
DEFAULT,2026-04-20T14:59:49.793038Z,METRIC_MEM:,1945.44,MB
34+
DEFAULT,2026-04-20T14:59:50.796118Z,METRIC_MEM:,2140.77,MB
35+
DEFAULT,2026-04-20T14:59:51.793792Z,METRIC_MEM:,2344.64,MB
36+
DEFAULT,2026-04-20T14:59:52.794297Z,METRIC_MEM:,2603.11,MB
37+
DEFAULT,2026-04-20T14:59:53.795767Z,METRIC_MEM:,2983.83,MB
38+
DEFAULT,2026-04-20T14:59:54.796083Z,METRIC_MEM:,3151.56,MB
39+
DEFAULT,2026-04-20T14:59:55.822665Z,METRIC_MEM:,3624.16,MB
40+
DEFAULT,2026-04-20T14:59:56.818814Z,METRIC_MEM:,3774.27,MB
41+
DEFAULT,2026-04-20T14:59:57.818682Z,METRIC_MEM:,4308.78,MB
42+
DEFAULT,2026-04-20T14:59:58.818511Z,METRIC_MEM:,4809.82,MB
43+
DEFAULT,2026-04-20T14:59:59.821162Z,METRIC_MEM:,5383.96,MB
44+
DEFAULT,2026-04-20T15:00:00.818704Z,METRIC_MEM:,5452.28,MB
45+
DEFAULT,2026-04-20T15:00:01.819104Z,METRIC_MEM:,5468.47,MB
46+
DEFAULT,2026-04-20T15:00:02.819475Z,METRIC_MEM:,5457.18,MB
47+
DEFAULT,2026-04-20T15:00:03.819830Z,METRIC_MEM:,5447.39,MB
48+
DEFAULT,2026-04-20T15:00:04.820259Z,METRIC_MEM:,5465.26,MB
49+
DEFAULT,2026-04-20T15:00:05.820598Z,METRIC_MEM:,5463.11,MB
50+
DEFAULT,2026-04-20T15:00:06.820974Z,METRIC_MEM:,5449.82,MB
51+
DEFAULT,2026-04-20T15:00:07.821433Z,METRIC_MEM:,5462.01,MB
52+
DEFAULT,2026-04-20T15:00:08.821801Z,METRIC_MEM:,5461.89,MB
53+
DEFAULT,2026-04-20T15:00:09.822611Z,METRIC_MEM:,5466.36,MB
54+
DEFAULT,2026-04-20T15:00:10.839283Z,METRIC_MEM:,5462.5,MB
55+
DEFAULT,2026-04-20T15:00:11.826381Z,METRIC_MEM:,5461.91,MB
56+
DEFAULT,2026-04-20T15:00:12.826815Z,METRIC_MEM:,5458.66,MB
57+
DEFAULT,2026-04-20T15:00:13.827094Z,METRIC_MEM:,5472.28,MB
58+
DEFAULT,2026-04-20T15:00:14.830802Z,METRIC_MEM:,5777.54,MB
59+
DEFAULT,2026-04-20T15:00:15.834612Z,METRIC_MEM:,6242.66,MB
60+
DEFAULT,2026-04-20T15:00:16.831243Z,METRIC_MEM:,6952.15,MB
61+
DEFAULT,2026-04-20T15:00:17.832603Z,METRIC_MEM:,7086.52,MB
62+
DEFAULT,2026-04-20T15:00:18.837562Z,METRIC_MEM:,7565.95,MB
63+
DEFAULT,2026-04-20T15:00:19.832365Z,METRIC_MEM:,7538.39,MB
64+
DEFAULT,2026-04-20T15:00:20.837168Z,METRIC_MEM:,7572.7,MB
65+
DEFAULT,2026-04-20T15:00:21.833118Z,METRIC_MEM:,7471.39,MB
66+
DEFAULT,2026-04-20T15:00:22.833394Z,METRIC_MEM:,7300.56,MB
67+
DEFAULT,2026-04-20T15:00:23.834075Z,METRIC_MEM:,7132.16,MB
68+
DEFAULT,2026-04-20T15:00:24.834825Z,METRIC_MEM:,6974.58,MB
69+
DEFAULT,2026-04-20T15:00:25.835584Z,METRIC_MEM:,6812.74,MB
70+
DEFAULT,2026-04-20T15:00:26.836323Z,METRIC_MEM:,6683.75,MB
71+
DEFAULT,2026-04-20T15:00:27.837011Z,METRIC_MEM:,6570.68,MB
72+
DEFAULT,2026-04-20T15:00:28.840331Z,METRIC_MEM:,6553.34,MB
73+
DEFAULT,2026-04-20T15:00:29.842342Z,METRIC_MEM:,6449.74,MB
74+
DEFAULT,2026-04-20T15:00:30.842779Z,METRIC_MEM:,6420.25,MB
75+
DEFAULT,2026-04-20T15:00:31.843102Z,METRIC_MEM:,6447.04,MB
76+
DEFAULT,2026-04-20T15:00:32.846716Z,METRIC_MEM:,6428.24,MB
77+
DEFAULT,2026-04-20T15:00:33.850433Z,METRIC_MEM:,7085.44,MB
78+
DEFAULT,2026-04-20T15:00:34.847560Z,METRIC_MEM:,7094.02,MB
79+
DEFAULT,2026-04-20T15:00:35.847891Z,METRIC_MEM:,7108.24,MB
80+
DEFAULT,2026-04-20T15:00:36.847936Z,METRIC_MEM:,7095.18,MB
81+
DEFAULT,2026-04-20T15:00:37.847913Z,METRIC_MEM:,7092.72,MB
82+
DEFAULT,2026-04-20T15:00:38.847807Z,METRIC_MEM:,7062.34,MB
83+
DEFAULT,2026-04-20T15:00:39.853210Z,METRIC_MEM:,6921.07,MB
84+
DEFAULT,2026-04-20T15:00:40.853684Z,METRIC_MEM:,6909.97,MB
85+
DEFAULT,2026-04-20T15:00:41.854214Z,METRIC_MEM:,6891.98,MB
86+
DEFAULT,2026-04-20T15:00:42.854490Z,METRIC_MEM:,6879.31,MB
87+
DEFAULT,2026-04-20T15:00:43.855020Z,METRIC_MEM:,6877.95,MB
88+
DEFAULT,2026-04-20T15:00:44.855450Z,METRIC_MEM:,6894.8,MB
89+
DEFAULT,2026-04-20T15:00:45.855873Z,METRIC_MEM:,6756.96,MB
90+
DEFAULT,2026-04-20T15:00:46.857784Z,METRIC_MEM:,6485.72,MB
91+
DEFAULT,2026-04-20T15:00:47.856652Z,METRIC_MEM:,6450.91,MB
92+
DEFAULT,2026-04-20T15:00:48.857031Z,METRIC_MEM:,6517.63,MB
93+
DEFAULT,2026-04-20T15:00:49.862250Z,METRIC_MEM:,6513.41,MB
94+
DEFAULT,2026-04-20T15:00:50.857773Z,METRIC_MEM:,6521.28,MB
95+
DEFAULT,2026-04-20T15:00:51.858133Z,METRIC_MEM:,6505.07,MB
96+
DEFAULT,2026-04-20T15:00:52.858414Z,METRIC_MEM:,6429.95,MB
97+
DEFAULT,2026-04-20T15:00:53.858797Z,METRIC_MEM:,6476.4,MB
98+
DEFAULT,2026-04-20T15:00:54.859127Z,METRIC_MEM:,6480.01,MB
99+
DEFAULT,2026-04-20T15:00:55.859485Z,METRIC_MEM:,6481.43,MB
100+
DEFAULT,2026-04-20T15:00:56.859441Z,METRIC_MEM:,6478.38,MB
101+
DEFAULT,2026-04-20T15:00:57.859322Z,METRIC_MEM:,6479.16,MB
102+
DEFAULT,2026-04-20T15:00:58.859179Z,METRIC_MEM:,6474.57,MB
103+
DEFAULT,2026-04-20T15:00:59.859043Z,METRIC_MEM:,6344.71,MB
104+
DEFAULT,2026-04-20T15:01:00.859298Z,METRIC_MEM:,6363.68,MB
105+
DEFAULT,2026-04-20T15:01:01.859673Z,METRIC_MEM:,6307.33,MB
106+
DEFAULT,2026-04-20T15:01:02.859970Z,METRIC_MEM:,6302.04,MB
107+
DEFAULT,2026-04-20T15:01:03.860228Z,METRIC_MEM:,6297.34,MB
108+
DEFAULT,2026-04-20T15:01:04.860579Z,METRIC_MEM:,6293.2,MB
109+
DEFAULT,2026-04-20T15:01:05.860960Z,METRIC_MEM:,6287.33,MB
110+
DEFAULT,2026-04-20T15:01:06.861270Z,METRIC_MEM:,6284.43,MB
111+
DEFAULT,2026-04-20T15:01:07.861650Z,METRIC_MEM:,6284.43,MB

assets/benchmarks/polars/README.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
This section details the methodology used to capture the memory metrics in the [`GCP Stress-Test Metrics (Scaling Efficiency)`](/README.md#gcp-stress-test-metrics-scaling-efficiency)
44

5-
The telemetry logger below was added **temporarily** to the orchestrator for a specific benchmarking run. This code was pushed directly to the Cloud Artifact Registry as an experimental image tag (`mem-record`) and is not part of the permanent git repository history.
5+
The telemetry logger below was added to the orchestrator for a specific benchmarking run.
66

77
```python
88
import psutil
@@ -28,7 +28,7 @@ finally:
2828
stop_event.set()
2929
logger_thread.join()
3030
```
31-
Since `psutil` requires C-extensions to compile, the **Dockerfile** was modified to include the necessary build tools and the package itself. This allowed for benchmarking without altering the project's permanent `requirements.txt`.
31+
Since `psutil` requires C-extensions to compile, the **Dockerfile** was modified to include the necessary build tools and the package itself. This allowed for benchmarking without altering the project's permanent [`requirements.txt`](/data_pipeline/requirements.txt).
3232

3333
```docker
3434
FROM python:3.11-slim

data_pipeline/Dockerfile

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
FROM python:3.11-slim
1+
FROM python:3.12-slim
22

33
ENV PYTHONDONTWRITEBYTECODE=1
44
ENV PYTHONUNBUFFERED=1

data_pipeline/assembly/assembly_executor.py

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,15 @@
33
# =============================================================================
44

55
import gc
6+
from typing import Dict
67
import ctypes
78
import platform
8-
from typing import Dict
99
from data_pipeline.shared.run_context import RunContext
10-
from data_pipeline.shared.loader_exporter import load_historical_data, export_file
10+
from data_pipeline.shared.loader_exporter import (
11+
load_historical_data,
12+
scan_gcs_uris_from_bigquery,
13+
export_file,
14+
)
1115
from data_pipeline.shared.modeling_configs import DIMENSION_REFERENCES
1216
from data_pipeline.assembly.assembly_logic import (
1317
init_report,
@@ -181,20 +185,26 @@ def orchestrate_dimension_refs(run_context: RunContext, report: Dict) -> bool:
181185
lf_raw = None
182186
df_dim = None
183187

184-
base_contracted_path = run_context.contracted_path
185-
186188
try:
187-
lf_raw = load_historical_data(
188-
base_path=base_contracted_path,
189-
table_name=table,
190-
log_info=lambda msg: loaded_data(msg, report),
191-
)
189+
# Switch between local and gcp IO
190+
if run_context.bq_project_id == "PROJECT_ID_NOT_DETECTED":
191+
lf_raw = load_historical_data(
192+
base_path=run_context.storage_contracted_path, table_name=table
193+
)
194+
else:
195+
lf_raw = scan_gcs_uris_from_bigquery(
196+
project_id=run_context.bq_project_id,
197+
dataset_id=run_context.bq_dataset_id,
198+
table_id=table,
199+
log_info=lambda msg: loaded_data(msg, report),
200+
)
192201

193202
if lf_raw is None:
194203
return False
195204

196205
primary_key = config.get("primary_key", [])
197206
require_col = config.get("required_column", [])
207+
dtypes = config.get("dtypes", {})
198208

199209
ok, df_dim = task_wrapper(
200210
report=report,
@@ -204,6 +214,7 @@ def orchestrate_dimension_refs(run_context: RunContext, report: Dict) -> bool:
204214
lf=lf_raw,
205215
primary_key=primary_key,
206216
req_column=require_col,
217+
dtypes=dtypes,
207218
)
208219

209220
if not ok:

data_pipeline/assembly/assembly_logic.py

Lines changed: 36 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,10 @@
66
from pathlib import Path
77
from typing import Dict, Callable, Any, List
88
from data_pipeline.shared.run_context import RunContext
9-
from data_pipeline.shared.loader_exporter import load_historical_data
9+
from data_pipeline.shared.loader_exporter import (
10+
load_historical_data,
11+
scan_gcs_uris_from_bigquery,
12+
)
1013
from data_pipeline.shared.modeling_configs import ASSEMBLE_SCHEMA, ASSEMBLE_DTYPES
1114

1215
EVENT_TABLES = ["df_orders", "df_order_items", "df_payments"]
@@ -151,7 +154,7 @@ def derive_fields(lf: pl.LazyFrame) -> pl.LazyFrame:
151154
)
152155
.dt.total_days()
153156
.cast(pl.Int16),
154-
order_date=pl.col("order_purchase_timestamp").dt.date(),
157+
order_date=pl.col("order_purchase_timestamp").dt.date().cast(pl.Datetime("us")),
155158
order_year_week=pl.col("order_purchase_timestamp")
156159
.dt.strftime("%G-W%V")
157160
.cast(pl.Categorical),
@@ -187,9 +190,17 @@ def freeze_schema(lf: pl.LazyFrame) -> pl.LazyFrame:
187190
if missing_cols:
188191
raise RuntimeError(f"missing required columns: {sorted(missing_cols)}")
189192

190-
lf_contract = lf.select(ASSEMBLE_SCHEMA).cast(pl.Schema(ASSEMBLE_DTYPES))
193+
lf_contract = lf.select(ASSEMBLE_SCHEMA)
191194

192-
return lf_contract
195+
datetime_cols = [
196+
col for col, dtype in ASSEMBLE_DTYPES.items() if isinstance(dtype, pl.Datetime)
197+
]
198+
199+
lf_contract = lf_contract.with_columns(
200+
[pl.col(col).dt.cast_time_unit("us") for col in datetime_cols]
201+
)
202+
203+
return lf_contract.cast(pl.Schema(ASSEMBLE_DTYPES))
193204

194205

195206
# ------------------------------------------------------------
@@ -201,12 +212,14 @@ def dimension_references(
201212
lf: pl.LazyFrame,
202213
primary_key: list[str],
203214
req_column: list[str],
215+
dtypes: dict,
204216
) -> pl.LazyFrame:
205217
"""
206218
Extracts a unique reference dataset from a historical source.
207219
208220
Contract:
209221
- Subtractive Filtering: Selects specified 'req_column' set and enforces uniqueness.
222+
- Type Enforcement: Casts columns to the formats defined in the provided 'dtypes' schema.
210223
211224
Invariants:
212225
- Dataset Grain: Strictly one row per 'primary_key'.
@@ -218,7 +231,7 @@ def dimension_references(
218231
- [Structural] Crashes if input LazyFrame lacks 'primary_key' or 'req_column'.
219232
"""
220233

221-
lf_dim = lf.select(req_column).unique(subset=primary_key)
234+
lf_dim = lf.select(req_column).unique(subset=primary_key).cast(pl.Schema(dtypes))
222235

223236
return lf_dim
224237

@@ -283,10 +296,10 @@ def task_wrapper(
283296

284297
def load_event_table(run_context: RunContext, report: Dict) -> Any:
285298
"""
286-
Batch-loads core event tables required for assembly.
299+
Batch-loads core event tables required for assembly from BigQuery.
287300
288301
Contract:
289-
- Hydrate: Iterates through EVENT_TABLES and loads Parquet files from 'contracted_path'.
302+
- Hydrate: Iterates through EVENT_TABLES and streams data via scan_gcs_uris_from_bigquery.
290303
291304
Outputs:
292305
- Dict keyed by table name containing loaded LazyFrames.
@@ -295,22 +308,31 @@ def load_event_table(run_context: RunContext, report: Dict) -> Any:
295308
- [Operational] Returns None if any required table is missing or fails to load.
296309
"""
297310

298-
base_contracted_path = run_context.contracted_path
299311
tables = {}
300312

301313
for table_name in EVENT_TABLES:
302314
try:
303-
df = load_historical_data(
304-
base_path=base_contracted_path,
305-
table_name=table_name,
306-
log_info=lambda msg: loaded_data(msg, report),
307-
)
315+
# Switch between local and gcp IO
316+
if run_context.bq_project_id == "PROJECT_ID_NOT_DETECTED":
317+
df = load_historical_data(
318+
base_path=run_context.storage_contracted_path,
319+
table_name=table_name,
320+
log_info=lambda msg: loaded_data(msg, report),
321+
)
322+
323+
else:
324+
df = scan_gcs_uris_from_bigquery(
325+
project_id=run_context.bq_project_id,
326+
dataset_id=run_context.bq_dataset_id,
327+
table_id=table_name,
328+
log_info=lambda msg: loaded_data(msg, report),
329+
)
308330

309331
if df is not None:
310332
tables[table_name] = df
311333

312334
except Exception as e:
313-
log_error(f"Required table {table_name} not found: {e}", report)
335+
log_error(f"Required table {table_name} not found : {e}", report)
314336
return None
315337

316338
if len(tables) < len(EVENT_TABLES):

data_pipeline/requirements.txt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,4 +2,5 @@ polars==1.39.0
22
pyarrow==19.0.0
33
google-cloud-storage
44
google-cloud-bigquery>=3.0.0
5+
google-cloud-bigquery-storage>=2.36.0
56
psutil==5.9.8

0 commit comments

Comments
 (0)