Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions docs/examples/script/README.md
Original file line number Diff line number Diff line change
@@ -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 `<example-name>/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).
23 changes: 23 additions & 0 deletions docs/examples/script/account_rollup/README.md
Original file line number Diff line number Diff line change
@@ -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.
13 changes: 13 additions & 0 deletions docs/examples/script/account_rollup/payload/config.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"sdkVersion": "6.0.4",
"entryPoint": "entrypoint.py",
"dataspace": "default",
"permissions": {
"read": {
"dlo": ["Account__dll"]
},
"write": {
"dlo": ["Account_Rollup__dll"]
}
}
}
62 changes: 62 additions & 0 deletions docs/examples/script/account_rollup/payload/entrypoint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
from pyspark.sql.functions import (
coalesce,
col,
lit,
sum as _sum,
)

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()
1 change: 1 addition & 0 deletions docs/examples/script/account_rollup/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
# Packages required for the custom code
9 changes: 9 additions & 0 deletions docs/examples/script/account_rollup/sample_data/account.csv
Original file line number Diff line number Diff line change
@@ -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
33 changes: 33 additions & 0 deletions docs/examples/script/json_explode/README.md
Original file line number Diff line number Diff line change
@@ -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 | |
13 changes: 13 additions & 0 deletions docs/examples/script/json_explode/payload/config.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"sdkVersion": "6.0.4",
"entryPoint": "entrypoint.py",
"dataspace": "default",
"permissions": {
"read": {
"dlo": ["User_Sessions__dll"]
},
"write": {
"dlo": ["User_Events__dll"]
}
}
}
22 changes: 22 additions & 0 deletions docs/examples/script/json_explode/payload/entrypoint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
from events import parse_events
from pyspark.sql.functions import col

from datacustomcode.client import Client
from datacustomcode.io.writer.base import WriteMode


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()
50 changes: 50 additions & 0 deletions docs/examples/script/json_explode/payload/py-files/events.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
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<Struct> 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"),
)
1 change: 1 addition & 0 deletions docs/examples/script/json_explode/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
# Packages required for the custom code
Original file line number Diff line number Diff line change
@@ -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}]"
33 changes: 33 additions & 0 deletions docs/examples/script/lead_scoring/README.md
Original file line number Diff line number Diff line change
@@ -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).
13 changes: 13 additions & 0 deletions docs/examples/script/lead_scoring/payload/config.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"sdkVersion": "6.0.4",
"entryPoint": "entrypoint.py",
"dataspace": "default",
"permissions": {
"read": {
"dlo": ["Lead__dll"]
},
"write": {
"dlo": ["Lead_Scored__dll"]
}
}
}
85 changes: 85 additions & 0 deletions docs/examples/script/lead_scoring/payload/entrypoint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
import itertools
from pathlib import Path
import tempfile
import zipfile

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()
4 changes: 4 additions & 0 deletions docs/examples/script/lead_scoring/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
# Packages required for the custom code
scikit-learn>=1.4
joblib>=1.3
pandas>=2.0
Loading
Loading