From a4a61d6048adba1bf854247067a58fe3624fe5ad Mon Sep 17 00:00:00 2001 From: Mark DeLaVergne Date: Wed, 23 Sep 2026 09:24:52 -0400 Subject: [PATCH 1/2] Add 3 new scripts into new docs/examples --- docs/examples/script/README.md | 16 +++ docs/examples/script/account_rollup/README.md | 23 ++++ .../script/account_rollup/payload/config.json | 13 ++ .../account_rollup/payload/entrypoint.py | 57 +++++++++ .../script/account_rollup/requirements.txt | 1 + .../account_rollup/sample_data/account.csv | 9 ++ docs/examples/script/json_explode/README.md | 33 +++++ .../script/json_explode/payload/config.json | 13 ++ .../script/json_explode/payload/entrypoint.py | 22 ++++ .../json_explode/payload/py-files/events.py | 41 ++++++ .../script/json_explode/requirements.txt | 1 + .../sample_data/user_sessions.csv | 6 + docs/examples/script/lead_scoring/README.md | 33 +++++ .../script/lead_scoring/payload/config.json | 13 ++ .../script/lead_scoring/payload/entrypoint.py | 82 ++++++++++++ .../script/lead_scoring/requirements.txt | 4 + .../script/lead_scoring/sample_data/lead.csv | 21 ++++ .../script/lead_scoring/train_model.py | 117 ++++++++++++++++++ 18 files changed, 505 insertions(+) create mode 100644 docs/examples/script/README.md create mode 100644 docs/examples/script/account_rollup/README.md create mode 100644 docs/examples/script/account_rollup/payload/config.json create mode 100644 docs/examples/script/account_rollup/payload/entrypoint.py create mode 100644 docs/examples/script/account_rollup/requirements.txt create mode 100644 docs/examples/script/account_rollup/sample_data/account.csv create mode 100644 docs/examples/script/json_explode/README.md create mode 100644 docs/examples/script/json_explode/payload/config.json create mode 100644 docs/examples/script/json_explode/payload/entrypoint.py create mode 100644 docs/examples/script/json_explode/payload/py-files/events.py create mode 100644 docs/examples/script/json_explode/requirements.txt create mode 100644 docs/examples/script/json_explode/sample_data/user_sessions.csv create mode 100644 docs/examples/script/lead_scoring/README.md create mode 100644 docs/examples/script/lead_scoring/payload/config.json create mode 100644 docs/examples/script/lead_scoring/payload/entrypoint.py create mode 100644 docs/examples/script/lead_scoring/requirements.txt create mode 100644 docs/examples/script/lead_scoring/sample_data/lead.csv create mode 100644 docs/examples/script/lead_scoring/train_model.py diff --git a/docs/examples/script/README.md b/docs/examples/script/README.md new file mode 100644 index 0000000..c0481cd --- /dev/null +++ b/docs/examples/script/README.md @@ -0,0 +1,16 @@ +# Code Extension for Batch Transform — Examples + +End-to-end Code Extension scripts. Each subdirectory is a self-contained deploy package (`payload/entrypoint.py` + `config.json` + `requirements.txt` + any bundled assets) that can run against a Data 360 org. + +## Loading sample data + +Many examples ship with a `sample_data/` folder containing CSV(s) needed to exercise it. To load one: + +1. In Data 360, go to **Data Streams → New** and pick the **File Upload** connector. +2. Upload the CSV from `/sample_data/`. +3. Map columns to the target object listed in the example's README, in your dataspace. +4. Run the ingest. + +## Related + +- Developer guide — [Data 360 Code Extension](https://developer.salesforce.com/docs/data/data-cloud-code-ext/guide/use-custom-code.html). diff --git a/docs/examples/script/account_rollup/README.md b/docs/examples/script/account_rollup/README.md new file mode 100644 index 0000000..be3c76e --- /dev/null +++ b/docs/examples/script/account_rollup/README.md @@ -0,0 +1,23 @@ +# Example: Iterative rollup of an account hierarchy + +A Code Extension for Batch Transform that walks a self-referencing `Account__dll` and produces `Account_Rollup__dll` with each account's **total-tree ARR** — its own ARR plus every descendant account's ARR. + +## What it demonstrates + +- Imperative control flow (a `while` loop with `isEmpty()` convergence) around DataFrame ops. +- `.persist()` on the frontier + growing set, so each iteration doesn't re-derive earlier hops. +- Alias-based joins (`accounts.alias("a")` / `descendants.alias("d")`) — required whenever the same DataFrame appears on both sides. + +## Prerequisites + +An `Account__dll` DLO in the target dataspace with columns: + +| Column | Type | Notes | +|-----------------|---------|---------------------------------------------| +| `id__c` | text | Primary key. | +| `parent_id__c` | text | Nullable. FK back to `Account__dll.id__c`. | +| `arr__c` | number | Annual recurring revenue for this account. | + +Data should be loaded via `sample_data/account.csv`. See [loading sample data](../README.md#loading-sample-data). + +Output DLO `Account_Rollup__dll` must also exist with `id__c` and `tree_arr__c` columns. diff --git a/docs/examples/script/account_rollup/payload/config.json b/docs/examples/script/account_rollup/payload/config.json new file mode 100644 index 0000000..a2f4973 --- /dev/null +++ b/docs/examples/script/account_rollup/payload/config.json @@ -0,0 +1,13 @@ +{ + "sdkVersion": "6.0.4", + "entryPoint": "entrypoint.py", + "dataspace": "default", + "permissions": { + "read": { + "dlo": ["Account__dll"] + }, + "write": { + "dlo": ["Account_Rollup__dll"] + } + } +} diff --git a/docs/examples/script/account_rollup/payload/entrypoint.py b/docs/examples/script/account_rollup/payload/entrypoint.py new file mode 100644 index 0000000..5eefd46 --- /dev/null +++ b/docs/examples/script/account_rollup/payload/entrypoint.py @@ -0,0 +1,57 @@ +from pyspark.sql.functions import col, sum as _sum, coalesce, lit + +from datacustomcode.client import Client +from datacustomcode.io.writer.base import WriteMode + + +def main(): + client = Client() + + accounts = ( + client.read_dlo("Account__dll") + .select("id__c", "parent_id__c", "arr__c") + .persist() + ) + + descendants = accounts.select( + col("id__c").alias("root_id"), + col("id__c").alias("descendant_id"), + ).persist() + + frontier = descendants + while True: + next_hop = ( + frontier.alias("f") + .join( + accounts.alias("a"), + col("a.parent_id__c") == col("f.descendant_id"), + "inner", + ) + .select( + col("f.root_id").alias("root_id"), + col("a.id__c").alias("descendant_id"), + ) + ) + if next_hop.isEmpty(): + break + descendants = descendants.union(next_hop).persist() + frontier = next_hop + + totals = ( + descendants.alias("d") + .join( + accounts.alias("a"), + col("d.descendant_id") == col("a.id__c"), + "left", + ) + .select(col("d.root_id"), col("a.arr__c")) + .groupBy("root_id") + .agg(coalesce(_sum("arr__c"), lit(0)).alias("tree_arr__c")) + .withColumnRenamed("root_id", "id__c") + ) + + client.write_to_dlo("Account_Rollup__dll", totals, WriteMode.OVERWRITE) + + +if __name__ == "__main__": + main() diff --git a/docs/examples/script/account_rollup/requirements.txt b/docs/examples/script/account_rollup/requirements.txt new file mode 100644 index 0000000..b4fc884 --- /dev/null +++ b/docs/examples/script/account_rollup/requirements.txt @@ -0,0 +1 @@ +# Packages required for the custom code diff --git a/docs/examples/script/account_rollup/sample_data/account.csv b/docs/examples/script/account_rollup/sample_data/account.csv new file mode 100644 index 0000000..c58cc42 --- /dev/null +++ b/docs/examples/script/account_rollup/sample_data/account.csv @@ -0,0 +1,9 @@ +id,parent_id,arr +A1,,100000 +A2,A1,50000 +A3,A1,25000 +A4,A2,10000 +A5,A2,5000 +A6,,200000 +A7,A6,75000 +A8,A7,15000 diff --git a/docs/examples/script/json_explode/README.md b/docs/examples/script/json_explode/README.md new file mode 100644 index 0000000..5f21bac --- /dev/null +++ b/docs/examples/script/json_explode/README.md @@ -0,0 +1,33 @@ +# Example: Flatten a batched event stream from a JSON-array column + +A Code Extension for Batch Transform that reads a `User_Sessions__dll` DLO — one row per session, with each session's events packed into a JSON-array column — and writes one row per event into `User_Events__dll`. Parse-and-explode logic lives in a reusable module under `payload/py-files/`. + +## What it demonstrates + +- **`payload/py-files/`** — Split real logic into modules, unit-test them locally, reuse them across scripts. +- **`from_json` + `explode_outer`** to turn one row per session into one row per event, carrying the parent `session_id__c` and `user_id__c` through. +- **`to_timestamp()`** to parse ISO-8601 event timestamps into Spark `TimestampType`, which writes cleanly into Data 360's `DateTime` column. + +## Prerequisites + +Input DLO `User_Sessions__dll` in the target dataspace: + +| Column | Type | Notes | +|------------------|------|---------------------------------------------| +| `session_id__c` | text | Primary key. | +| `user_id__c` | text | | +| `events__c` | text | JSON array of `{event_id, event_type, ts, path, value}`. | + +Data should be loaded via `sample_data/user_sessions.csv`. See [loading sample data](../README.md#loading-sample-data). + +Output DLO `User_Events__dll` in the same dataspace: + +| Column | Type | Notes | +|-----------------|----------|---------------| +| `event_id__c` | text | Primary key. | +| `session_id__c` | text | | +| `user_id__c` | text | | +| `event_type__c` | text | | +| `event_ts__c` | DateTime | | +| `path__c` | text | | +| `value__c` | text | | diff --git a/docs/examples/script/json_explode/payload/config.json b/docs/examples/script/json_explode/payload/config.json new file mode 100644 index 0000000..10a7c11 --- /dev/null +++ b/docs/examples/script/json_explode/payload/config.json @@ -0,0 +1,13 @@ +{ + "sdkVersion": "6.0.4", + "entryPoint": "entrypoint.py", + "dataspace": "default", + "permissions": { + "read": { + "dlo": ["User_Sessions__dll"] + }, + "write": { + "dlo": ["User_Events__dll"] + } + } +} diff --git a/docs/examples/script/json_explode/payload/entrypoint.py b/docs/examples/script/json_explode/payload/entrypoint.py new file mode 100644 index 0000000..ffa3be7 --- /dev/null +++ b/docs/examples/script/json_explode/payload/entrypoint.py @@ -0,0 +1,22 @@ +from pyspark.sql.functions import col + +from datacustomcode.client import Client +from datacustomcode.io.writer.base import WriteMode +from events import parse_events + + +def main(): + client = Client() + + sessions = client.read_dlo("User_Sessions__dll").select( + "session_id__c", "user_id__c", "events__c" + ) + exploded = parse_events(sessions, "session_id__c", "user_id__c", "events__c") + + output = exploded.filter(col("event_id__c").isNotNull()) + + client.write_to_dlo("User_Events__dll", output, WriteMode.OVERWRITE) + + +if __name__ == "__main__": + main() diff --git a/docs/examples/script/json_explode/payload/py-files/events.py b/docs/examples/script/json_explode/payload/py-files/events.py new file mode 100644 index 0000000..67d2a71 --- /dev/null +++ b/docs/examples/script/json_explode/payload/py-files/events.py @@ -0,0 +1,41 @@ +from pyspark.sql.functions import col, explode_outer, from_json, to_timestamp +from pyspark.sql.types import ( + ArrayType, StringType, StructField, StructType, +) + + +EVENT_SCHEMA = StructType([ + StructField("event_id", StringType(), True), + StructField("event_type", StringType(), True), + StructField("ts", StringType(), True), + StructField("path", StringType(), True), + StructField("value", StringType(), True), +]) +EVENT_ARRAY_SCHEMA = ArrayType(EVENT_SCHEMA, True) + + +def parse_events(df, session_col, user_col, events_col): + """One row per event, carrying its parent session + user. + + Handles both a real ArrayType column and a string column that + contains the JSON array. + """ + field = next(f for f in df.schema.fields if f.name == events_col) + events = ( + col(events_col) + if isinstance(field.dataType, ArrayType) + else from_json(col(events_col).cast("string"), EVENT_ARRAY_SCHEMA) + ) + return df.select( + col(session_col), + col(user_col), + explode_outer(events).alias("evt"), + ).select( + col("evt.event_id").alias("event_id__c"), + col(session_col), + col(user_col), + col("evt.event_type").alias("event_type__c"), + to_timestamp(col("evt.ts")).alias("event_ts__c"), + col("evt.path").alias("path__c"), + col("evt.value").alias("value__c"), + ) diff --git a/docs/examples/script/json_explode/requirements.txt b/docs/examples/script/json_explode/requirements.txt new file mode 100644 index 0000000..b4fc884 --- /dev/null +++ b/docs/examples/script/json_explode/requirements.txt @@ -0,0 +1 @@ +# Packages required for the custom code diff --git a/docs/examples/script/json_explode/sample_data/user_sessions.csv b/docs/examples/script/json_explode/sample_data/user_sessions.csv new file mode 100644 index 0000000..e09df3e --- /dev/null +++ b/docs/examples/script/json_explode/sample_data/user_sessions.csv @@ -0,0 +1,6 @@ +session_id,user_id,events +S1,U1,"[{""event_id"":""E1"",""event_type"":""page_view"",""ts"":""2026-09-22T14:00:00Z"",""path"":""/home"",""value"":null},{""event_id"":""E2"",""event_type"":""page_view"",""ts"":""2026-09-22T14:00:15Z"",""path"":""/pricing"",""value"":null},{""event_id"":""E3"",""event_type"":""click"",""ts"":""2026-09-22T14:00:42Z"",""path"":""/pricing"",""value"":""cta_upgrade""}]" +S2,U2,"[{""event_id"":""E4"",""event_type"":""page_view"",""ts"":""2026-09-22T14:05:00Z"",""path"":""/blog/announcing-x"",""value"":null},{""event_id"":""E5"",""event_type"":""scroll_50"",""ts"":""2026-09-22T14:05:22Z"",""path"":""/blog/announcing-x"",""value"":""50""}]" +S3,U1,"[{""event_id"":""E6"",""event_type"":""page_view"",""ts"":""2026-09-22T15:00:00Z"",""path"":""/dashboard"",""value"":null},{""event_id"":""E7"",""event_type"":""form_submit"",""ts"":""2026-09-22T15:01:11Z"",""path"":""/dashboard"",""value"":""support_request""},{""event_id"":""E8"",""event_type"":""page_view"",""ts"":""2026-09-22T15:01:20Z"",""path"":""/support/thanks"",""value"":null}]" +S4,U3,"[]" +S5,U4,"[{""event_id"":""E9"",""event_type"":""page_view"",""ts"":""2026-09-22T16:10:00Z"",""path"":""/home"",""value"":null}]" diff --git a/docs/examples/script/lead_scoring/README.md b/docs/examples/script/lead_scoring/README.md new file mode 100644 index 0000000..254d136 --- /dev/null +++ b/docs/examples/script/lead_scoring/README.md @@ -0,0 +1,33 @@ +# Example: Score leads with a bundled scikit-learn model + +A Code Extension for Batch Transform that scores every lead in a `Lead__dll` DLO against a scikit-learn classifier **bundled directly with the deploy package**. + +## What it demonstrates + +- **Bundling a trained model as `payload/files/lead_scorer.zip`** (a zip wrapping `lead_scorer.joblib`). +- **Grid-join scoring** — instead of scoring each lead individually, the script enumerates every combination of the model's categorical inputs, scores that small grid with `predict_proba`, and turns the results into a reference DataFrame that is left-joined against the leads. + +## Regenerating the bundled model + +```sh +pip install -r requirements.txt +python train_model.py +``` + +## Prerequisites + +A `Lead__dll` DLO in the target dataspace with at least: + +| Column | Type | +|------------------|--------| +| `id__c` | text | +| `first_name__c` | text | +| `last_name__c` | text | +| `industry__c` | text | +| `employee_band__c` | text | +| `region__c` | text | +| `source__c` | text | + +Data should be loaded via `sample_data/lead.csv`. See [loading sample data](../README.md#loading-sample-data). + +Output DLO `Lead_Scored__dll` must also exist (same columns plus `score__c` number). diff --git a/docs/examples/script/lead_scoring/payload/config.json b/docs/examples/script/lead_scoring/payload/config.json new file mode 100644 index 0000000..e05a595 --- /dev/null +++ b/docs/examples/script/lead_scoring/payload/config.json @@ -0,0 +1,13 @@ +{ + "sdkVersion": "6.0.4", + "entryPoint": "entrypoint.py", + "dataspace": "default", + "permissions": { + "read": { + "dlo": ["Lead__dll"] + }, + "write": { + "dlo": ["Lead_Scored__dll"] + } + } +} diff --git a/docs/examples/script/lead_scoring/payload/entrypoint.py b/docs/examples/script/lead_scoring/payload/entrypoint.py new file mode 100644 index 0000000..72420f2 --- /dev/null +++ b/docs/examples/script/lead_scoring/payload/entrypoint.py @@ -0,0 +1,82 @@ +import itertools +import tempfile +import zipfile +from pathlib import Path + +import joblib +import pandas as pd +from pyspark.sql import Row +from pyspark.sql.functions import coalesce, col, lit + +from datacustomcode.client import Client +from datacustomcode.io.writer.base import WriteMode + + +UNKNOWN = "__unknown__" +MODEL_ARCHIVE = "lead_scorer.zip" +MODEL_MEMBER = "lead_scorer.joblib" + +OUTPUT_COLUMNS = [ + "id__c", + "first_name__c", + "last_name__c", + "industry__c", + "employee_band__c", + "region__c", + "source__c", + "score__c", +] + + +def load_bundled_model(client): + archive_path = client.find_file_path(MODEL_ARCHIVE) + extract_dir = Path(tempfile.mkdtemp(prefix="lead_scorer_")) + with zipfile.ZipFile(archive_path) as zf: + zf.extract(MODEL_MEMBER, path=extract_dir) + return joblib.load(extract_dir / MODEL_MEMBER) + + +def build_score_grid(spark, bundle): + """Enumerate every combination of the model's known feature values, + score them, and return a small Spark DataFrame.""" + pipeline = bundle["pipeline"] + feature_cols = bundle["feature_cols"] + domain = bundle["feature_domain"] + + combos = list(itertools.product(*(domain[c] for c in feature_cols))) + grid_pd = pd.DataFrame(combos, columns=feature_cols) + grid_pd["score__c"] = pipeline.predict_proba(grid_pd)[:, 1] + + rows = [ + Row(**{c: r[c] for c in feature_cols}, score__c=float(r["score__c"])) + for r in grid_pd.to_dict(orient="records") + ] + return spark.createDataFrame(rows), feature_cols + + +def main(): + client = Client() + + leads = client.read_dlo("Lead__dll") + spark = leads.sparkSession + + bundle = load_bundled_model(client) + grid, feature_cols = build_score_grid(spark, bundle) + + normalized = leads + for c in feature_cols: + normalized = normalized.withColumn(c, coalesce(col(c), lit(UNKNOWN))) + + scored = ( + normalized.alias("l") + .join(grid.alias("g"), feature_cols, "left") + .withColumn("score__c", coalesce(col("score__c"), lit(0.0))) + ) + + client.write_to_dlo( + "Lead_Scored__dll", scored.select(*OUTPUT_COLUMNS), WriteMode.OVERWRITE + ) + + +if __name__ == "__main__": + main() diff --git a/docs/examples/script/lead_scoring/requirements.txt b/docs/examples/script/lead_scoring/requirements.txt new file mode 100644 index 0000000..7356fa5 --- /dev/null +++ b/docs/examples/script/lead_scoring/requirements.txt @@ -0,0 +1,4 @@ +# Packages required for the custom code +scikit-learn>=1.4 +joblib>=1.3 +pandas>=2.0 diff --git a/docs/examples/script/lead_scoring/sample_data/lead.csv b/docs/examples/script/lead_scoring/sample_data/lead.csv new file mode 100644 index 0000000..70ec8af --- /dev/null +++ b/docs/examples/script/lead_scoring/sample_data/lead.csv @@ -0,0 +1,21 @@ +id,first_name,last_name,industry,employee_band,region,source +L1,Ava,Nguyen,tech,201-1000,NA,referral +L2,Ben,Ortiz,tech,1001-5000,NA,partner +L3,Cara,Singh,finance,5001+,EMEA,event +L4,Dan,Weiss,healthcare,51-200,NA,marketing_campaign +L5,Elena,Park,retail,11-50,APAC,cold_outreach +L6,Frank,Yamada,agriculture,1-10,LATAM,cold_outreach +L7,Grace,Osei,tech,51-200,EMEA,webform +L8,Hugo,Lin,telecom,1001-5000,ANZ,partner +L9,Isha,Mehta,education,1-10,APAC,marketing_campaign +L10,Jack,Rowe,government,201-1000,NA,cold_outreach +L11,Kira,Sato,media,11-50,APAC,webform +L12,Leo,Duval,real_estate,51-200,EMEA,referral +L13,Mia,Costa,hospitality,1-10,LATAM,webform +L14,Nate,Weber,energy,5001+,NA,partner +L15,Omar,Diallo,transportation,201-1000,EMEA,event +L16,Priya,Rao,manufacturing,1001-5000,APAC,referral +L17,Quinn,Bauer,other,11-50,ANZ,cold_outreach +L18,Rita,Kang,finance,51-200,APAC,referral +L19,Sam,Alavi,tech,1-10,NA,event +L20,Tara,Ilic,healthcare,5001+,EMEA,partner diff --git a/docs/examples/script/lead_scoring/train_model.py b/docs/examples/script/lead_scoring/train_model.py new file mode 100644 index 0000000..03a20da --- /dev/null +++ b/docs/examples/script/lead_scoring/train_model.py @@ -0,0 +1,117 @@ +# Train a small lead-scoring classifier and bundle it with the payload. +from __future__ import annotations + +import itertools +import pathlib +import random +import zipfile + +import joblib +import numpy as np +import pandas as pd +from sklearn.linear_model import LogisticRegression +from sklearn.metrics import roc_auc_score +from sklearn.model_selection import train_test_split +from sklearn.pipeline import Pipeline +from sklearn.preprocessing import OneHotEncoder + + +INDUSTRIES = [ + "retail", "healthcare", "finance", "tech", "manufacturing", + "education", "government", "media", "hospitality", "transportation", + "energy", "real_estate", "telecom", "agriculture", "other", +] +EMPLOYEE_BANDS = ["1-10", "11-50", "51-200", "201-1000", "1001-5000", "5001+"] +REGIONS = ["NA", "EMEA", "APAC", "LATAM", "ANZ"] +SOURCES = [ + "webform", "event", "referral", "partner", + "cold_outreach", "marketing_campaign", +] + +FEATURE_COLS = ["industry__c", "employee_band__c", "region__c", "source__c"] + +BAND_SIGNAL = {b: i for i, b in enumerate(EMPLOYEE_BANDS)} +SOURCE_SIGNAL = { + "referral": 2.0, "partner": 1.2, "event": 0.8, + "marketing_campaign": 0.3, "webform": 0.0, "cold_outreach": -1.0, +} +INDUSTRY_SIGNAL = { + "tech": 1.2, "finance": 1.0, "healthcare": 0.7, + "telecom": 0.5, "energy": 0.4, "manufacturing": 0.2, + "retail": 0.0, "media": 0.0, "real_estate": -0.2, + "education": -0.3, "government": -0.4, "hospitality": -0.5, + "transportation": -0.5, "agriculture": -0.7, "other": -0.5, +} +REGION_SIGNAL = {"NA": 0.6, "EMEA": 0.4, "APAC": 0.3, "ANZ": 0.2, "LATAM": 0.0} + + +def synth_rows(n: int, seed: int = 42) -> pd.DataFrame: + rng = random.Random(seed) + rows = [] + for _ in range(n): + row = { + "industry__c": rng.choice(INDUSTRIES), + "employee_band__c": rng.choice(EMPLOYEE_BANDS), + "region__c": rng.choice(REGIONS), + "source__c": rng.choice(SOURCES), + } + logit = ( + INDUSTRY_SIGNAL[row["industry__c"]] + + 0.4 * BAND_SIGNAL[row["employee_band__c"]] + + REGION_SIGNAL[row["region__c"]] + + SOURCE_SIGNAL[row["source__c"]] + - 1.5 + ) + p = 1.0 / (1.0 + np.exp(-logit)) + row["converted"] = int(rng.random() < p) + rows.append(row) + return pd.DataFrame(rows) + + +def main() -> None: + df = synth_rows(8000) + X, y = df[FEATURE_COLS], df["converted"] + + X_tr, X_te, y_tr, y_te = train_test_split( + X, y, test_size=0.25, random_state=0, stratify=y + ) + pipe = Pipeline( + [ + ("ohe", OneHotEncoder(handle_unknown="ignore", sparse_output=False)), + ("clf", LogisticRegression(max_iter=1000, random_state=0)), + ] + ) + pipe.fit(X_tr, y_tr) + + auc = roc_auc_score(y_te, pipe.predict_proba(X_te)[:, 1]) + prevalence = float(y.mean()) + + grid_rows = list(itertools.product(INDUSTRIES, EMPLOYEE_BANDS, REGIONS, SOURCES)) + print(f"train AUC: {auc:.3f} | prevalence: {prevalence:.3f} | grid: {len(grid_rows)}") + + payload = { + "pipeline": pipe, + "feature_cols": FEATURE_COLS, + "feature_domain": { + "industry__c": INDUSTRIES, + "employee_band__c": EMPLOYEE_BANDS, + "region__c": REGIONS, + "source__c": SOURCES, + }, + } + + out_joblib = pathlib.Path(__file__).parent / "payload" / "files" / "lead_scorer.joblib" + out_joblib.parent.mkdir(parents=True, exist_ok=True) + joblib.dump(payload, out_joblib, compress=3) + print(f"wrote {out_joblib} ({out_joblib.stat().st_size} bytes)") + + out_zip = out_joblib.with_suffix(".zip") + with zipfile.ZipFile(out_zip, "w", zipfile.ZIP_DEFLATED) as zf: + zf.write(out_joblib, arcname=out_joblib.name) + print(f"wrote {out_zip} ({out_zip.stat().st_size} bytes)") + + out_joblib.unlink() + + +if __name__ == "__main__": + main() From 7130e6f58325f7a6423c3bc24422cf77158483d6 Mon Sep 17 00:00:00 2001 From: Mark DeLaVergne Date: Wed, 23 Sep 2026 09:33:58 -0400 Subject: [PATCH 2/2] Lint fixes --- .../account_rollup/payload/entrypoint.py | 7 ++- .../script/json_explode/payload/entrypoint.py | 2 +- .../json_explode/payload/py-files/events.py | 29 ++++++--- .../script/lead_scoring/payload/entrypoint.py | 9 ++- .../script/lead_scoring/train_model.py | 63 ++++++++++++++----- pyproject.toml | 12 ++-- 6 files changed, 86 insertions(+), 36 deletions(-) diff --git a/docs/examples/script/account_rollup/payload/entrypoint.py b/docs/examples/script/account_rollup/payload/entrypoint.py index 5eefd46..613cbe8 100644 --- a/docs/examples/script/account_rollup/payload/entrypoint.py +++ b/docs/examples/script/account_rollup/payload/entrypoint.py @@ -1,4 +1,9 @@ -from pyspark.sql.functions import col, sum as _sum, coalesce, lit +from pyspark.sql.functions import ( + coalesce, + col, + lit, + sum as _sum, +) from datacustomcode.client import Client from datacustomcode.io.writer.base import WriteMode diff --git a/docs/examples/script/json_explode/payload/entrypoint.py b/docs/examples/script/json_explode/payload/entrypoint.py index ffa3be7..aaa9825 100644 --- a/docs/examples/script/json_explode/payload/entrypoint.py +++ b/docs/examples/script/json_explode/payload/entrypoint.py @@ -1,8 +1,8 @@ +from events import parse_events from pyspark.sql.functions import col from datacustomcode.client import Client from datacustomcode.io.writer.base import WriteMode -from events import parse_events def main(): diff --git a/docs/examples/script/json_explode/payload/py-files/events.py b/docs/examples/script/json_explode/payload/py-files/events.py index 67d2a71..877ef0f 100644 --- a/docs/examples/script/json_explode/payload/py-files/events.py +++ b/docs/examples/script/json_explode/payload/py-files/events.py @@ -1,16 +1,25 @@ -from pyspark.sql.functions import col, explode_outer, from_json, to_timestamp +from pyspark.sql.functions import ( + col, + explode_outer, + from_json, + to_timestamp, +) from pyspark.sql.types import ( - ArrayType, StringType, StructField, StructType, + ArrayType, + StringType, + StructField, + StructType, ) - -EVENT_SCHEMA = StructType([ - StructField("event_id", StringType(), True), - StructField("event_type", StringType(), True), - StructField("ts", StringType(), True), - StructField("path", StringType(), True), - StructField("value", StringType(), True), -]) +EVENT_SCHEMA = StructType( + [ + StructField("event_id", StringType(), True), + StructField("event_type", StringType(), True), + StructField("ts", StringType(), True), + StructField("path", StringType(), True), + StructField("value", StringType(), True), + ] +) EVENT_ARRAY_SCHEMA = ArrayType(EVENT_SCHEMA, True) diff --git a/docs/examples/script/lead_scoring/payload/entrypoint.py b/docs/examples/script/lead_scoring/payload/entrypoint.py index 72420f2..1befa73 100644 --- a/docs/examples/script/lead_scoring/payload/entrypoint.py +++ b/docs/examples/script/lead_scoring/payload/entrypoint.py @@ -1,17 +1,20 @@ import itertools +from pathlib import Path import tempfile import zipfile -from pathlib import Path import joblib import pandas as pd from pyspark.sql import Row -from pyspark.sql.functions import coalesce, col, lit +from pyspark.sql.functions import ( + coalesce, + col, + lit, +) from datacustomcode.client import Client from datacustomcode.io.writer.base import WriteMode - UNKNOWN = "__unknown__" MODEL_ARCHIVE = "lead_scorer.zip" MODEL_MEMBER = "lead_scorer.joblib" diff --git a/docs/examples/script/lead_scoring/train_model.py b/docs/examples/script/lead_scoring/train_model.py index 03a20da..f942173 100644 --- a/docs/examples/script/lead_scoring/train_model.py +++ b/docs/examples/script/lead_scoring/train_model.py @@ -15,32 +15,61 @@ from sklearn.pipeline import Pipeline from sklearn.preprocessing import OneHotEncoder - INDUSTRIES = [ - "retail", "healthcare", "finance", "tech", "manufacturing", - "education", "government", "media", "hospitality", "transportation", - "energy", "real_estate", "telecom", "agriculture", "other", + "retail", + "healthcare", + "finance", + "tech", + "manufacturing", + "education", + "government", + "media", + "hospitality", + "transportation", + "energy", + "real_estate", + "telecom", + "agriculture", + "other", ] EMPLOYEE_BANDS = ["1-10", "11-50", "51-200", "201-1000", "1001-5000", "5001+"] REGIONS = ["NA", "EMEA", "APAC", "LATAM", "ANZ"] SOURCES = [ - "webform", "event", "referral", "partner", - "cold_outreach", "marketing_campaign", + "webform", + "event", + "referral", + "partner", + "cold_outreach", + "marketing_campaign", ] FEATURE_COLS = ["industry__c", "employee_band__c", "region__c", "source__c"] BAND_SIGNAL = {b: i for i, b in enumerate(EMPLOYEE_BANDS)} SOURCE_SIGNAL = { - "referral": 2.0, "partner": 1.2, "event": 0.8, - "marketing_campaign": 0.3, "webform": 0.0, "cold_outreach": -1.0, + "referral": 2.0, + "partner": 1.2, + "event": 0.8, + "marketing_campaign": 0.3, + "webform": 0.0, + "cold_outreach": -1.0, } INDUSTRY_SIGNAL = { - "tech": 1.2, "finance": 1.0, "healthcare": 0.7, - "telecom": 0.5, "energy": 0.4, "manufacturing": 0.2, - "retail": 0.0, "media": 0.0, "real_estate": -0.2, - "education": -0.3, "government": -0.4, "hospitality": -0.5, - "transportation": -0.5, "agriculture": -0.7, "other": -0.5, + "tech": 1.2, + "finance": 1.0, + "healthcare": 0.7, + "telecom": 0.5, + "energy": 0.4, + "manufacturing": 0.2, + "retail": 0.0, + "media": 0.0, + "real_estate": -0.2, + "education": -0.3, + "government": -0.4, + "hospitality": -0.5, + "transportation": -0.5, + "agriculture": -0.7, + "other": -0.5, } REGION_SIGNAL = {"NA": 0.6, "EMEA": 0.4, "APAC": 0.3, "ANZ": 0.2, "LATAM": 0.0} @@ -87,7 +116,9 @@ def main() -> None: prevalence = float(y.mean()) grid_rows = list(itertools.product(INDUSTRIES, EMPLOYEE_BANDS, REGIONS, SOURCES)) - print(f"train AUC: {auc:.3f} | prevalence: {prevalence:.3f} | grid: {len(grid_rows)}") + print( + f"train AUC: {auc:.3f} | prevalence: {prevalence:.3f} | grid: {len(grid_rows)}" + ) payload = { "pipeline": pipe, @@ -100,7 +131,9 @@ def main() -> None: }, } - out_joblib = pathlib.Path(__file__).parent / "payload" / "files" / "lead_scorer.joblib" + out_joblib = ( + pathlib.Path(__file__).parent / "payload" / "files" / "lead_scorer.joblib" + ) out_joblib.parent.mkdir(parents=True, exist_ok=True) joblib.dump(payload, out_joblib, compress=3) print(f"wrote {out_joblib} ({out_joblib.stat().st_size} bytes)") diff --git a/pyproject.toml b/pyproject.toml index 2d6938b..c0e868b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -87,10 +87,10 @@ warn_unused_configs = true [tool.poetry] include = [ - {path = "src/datacustomcode/templates/**/*"}, - {path = "src/datacustomcode/config.yaml"} + { path = "src/datacustomcode/templates/**/*" }, + { path = "src/datacustomcode/config.yaml" } ] -packages = [{include = "datacustomcode", from = "src"}] +packages = [{ include = "datacustomcode", from = "src" }] version = "0.0.0" [tool.poetry.build] @@ -115,7 +115,7 @@ coverage = ">=7.0.0,<8.0.0" ipykernel = "^6.29.5" mypy = "*" pipreqs = "*" -poetry-dynamic-versioning = {extras = ["plugin"], version = "^1.8.2"} +poetry-dynamic-versioning = { extras = ["plugin"], version = "^1.8.2" } pre-commit = "*" pytest = "*" pytest-cov = "*" @@ -127,7 +127,7 @@ types-requests = "*" [tool.poetry.plugins."poetry.plugin"] [tool.poetry.requires-plugins] -poetry-dynamic-versioning = {version = ">=1.0.0,<2.0.0", extras = ["plugin"]} +poetry-dynamic-versioning = { version = ">=1.0.0,<2.0.0", extras = ["plugin"] } [tool.poetry.scripts] datacustomcode = "datacustomcode.cli:cli" @@ -152,7 +152,6 @@ testpaths = ["tests"] [tool.ruff] fix = true line-length = 88 -target-version = 'py310' lint.ignore = [ # do not assign a lambda expression, use a def 'E731', @@ -211,6 +210,7 @@ lint.select = [ 'RUF', 'S102' ] +target-version = 'py310' [tool.setuptools.package_data] datacustomcode = ["templates/**/*", "config.yaml"]