Salesforce Developers Blog

Advanced Code Extension PySpark Examples for Batch Data Transforms

Avatar for mdelavergnemdelavergne
Code Extension scripts for batch data transform allow you to use custom Python logic to transform your data lake objects (DLOs) and data model objects (DMOs) in an isolated compute environment in Data 360. In this blog, we’ll cover advanced batch transform script examples.
Advanced Code Extension PySpark Examples for Batch Data Transforms
October 06, 2026

Code Extension scripts for batch data transform allow you to use custom Python logic to transform your data lake objects (DLOs) and data model objects (DMOs) in an isolated compute environment in Data 360. In this blog, we’ll cover advanced batch transform script examples.

Batch data transforms

Data 360’s native batch data transform feature allows for data preparation over DLOs and DMOs, including through operations like joins, filters, and aggregates. You can use the Transform Builder UI to create batch transforms on the platform. The Batch Data Transforms in Data 360 Trailhead module is a great starting point.

Code Extension 

Code Extension is a Data 360 feature that runs custom PySpark logic against DLOs and DMOs inside an isolated Spark compute environment. It powers multiple pipelines, including search-index chunking and batch transform logic, the use case that this post covers.

If you are not familiar with the basics of Code Extension, you can check out the following resources before going deeper into this blog:

Under the hood: the development workflow for Code Extension batch data transform scripts

With Code Extensions, you author a PySpark script that drives the same batch transform pipeline as the one used by native batch transforms. A batch transform Code Extension script runs under strict isolation, using the datacustomcode SDK‘s Client() for reading and writing to DLOs or DMOs.

Here’s a high-level overview of what happens on your (the developer’s) machine and what happens in Data 360:

Architecture diagram showing two zones - developer machine and Data 360 - and the commands, APIs, Data 360 objects, and compute used within both.

On the developer machine, a Python virtual environment is used for local execution involving the SDK, your custom PySpark script, and its dependencies. To run the extension locally, you use the following Salesforce Command Line Interface (SF CLI) for Code Extension command: sf data-code-extension script run. The command executes a Spark job on your computer, which reads DLOs or DMOs using the Connect API (by default, the first 1000 rows), and writes only to console output. This is where you’ll iterate until the logic works as intended.

The next step is to deploy to your org. To do so, use the SF CLI command: sf data-code-extension script deploy. The command zips up the project using Docker and uploads it via the Connect API into a Data 360 Transform and Code Extension pair.

When running the transform (now or on a schedule), Data 360 manages the Spark job, which executes your PySpark script to read from and write to DLOs or DMOs.

Local run (script run) Deployed run (script deploy + transform)
Runs on Your machine, in a Python virtual environment Data 360 compute
Reads default limit By default, limit to 1000. See docs for how to override. No default limit
Writes to Console output only Target DLO or DMO
Spark managed by The SF CLI on your computer Data 360

We recently published Run Complex Data Transformations in Data 360 with Code Extension, which gives you more information about the basics. Below, we build on that introduction and walk through three more examples: rolling up an account hierarchy, exploding a JSON event stream, and scoring leads with a bundled scikit-learn model.

Example 1: Iterative rollup of an account hierarchy

The use case: Roll up each account’s total-tree Annual Recurring Revenue (ARR), including its own plus every descendant’s, at any depth, into Account_Rollup__dll.

Diagram of a sample Account__dll tree (A1→A2→A4/A5, A1→A3, A6→A7→A8), where each account has an ARR value that needs to be rolled up

The solution: With the Account__dll records having various descendant depths, we can use a loop and keep joining the frontier one hop deeper until nothing new comes back.

1from pyspark.sql.functions import col, sum as _sum, coalesce, lit
2
3from datacustomcode.client import Client
4from datacustomcode.io.writer.base import WriteMode
5
6
7def main():
8    client = Client()
9
10    accounts = (
11        client.read_dlo("Account__dll")
12        .select("id__c", "parent_id__c", "arr__c")
13        .persist()
14    )
15
16    descendants = accounts.select(
17        col("id__c").alias("root_id"),
18        col("id__c").alias("descendant_id"),
19    ).persist()
20
21    frontier = descendants
22    while True:
23        next_hop = (
24            frontier.alias("f")
25            .join(
26                accounts.alias("a"),
27                col("a.parent_id__c") == col("f.descendant_id"),
28                "inner",
29            )
30            .select(
31                col("f.root_id").alias("root_id"),
32                col("a.id__c").alias("descendant_id"),
33            )
34        )
35        if next_hop.isEmpty():
36            break
37        descendants = descendants.union(next_hop).persist()
38        frontier = next_hop
39
40    totals = (
41        descendants.alias("d")
42        .join(
43            accounts.alias("a"),
44            col("d.descendant_id") == col("a.id__c"),
45            "left",
46        )
47        .select(col("d.root_id"), col("a.arr__c"))
48        .groupBy("root_id")
49        .agg(coalesce(_sum("arr__c"), lit(0)).alias("tree_arr__c"))
50        .withColumnRenamed("root_id", "id__c")
51    )
52
53    client.write_to_dlo("Account_Rollup__dll", totals, WriteMode.OVERWRITE)
54
55
56if __name__ == "__main__":
57    main()

The convergence check, stop, accumulate loop is plain Python driving DataFrames. NOTE: This convergence check assumes an acyclic hierarchy. If parent_id__c can form a cycle, add a max-depth counter to break the loop so a bad record can’t hang the Spark job indefinitely.

The results written to Account_Rollup__dll are as expected for this sample dataset:

Query results showing Account_Rollup__dll records with correct, rolled up values for the tree_arr__c column.

You can find this entire example, including the sample data used, in the SDK’s docs/examples/script/account_rollup directory.

Example 2: Flatten a batched event stream from a JSON-array column

The use case: Analytics SDKs commonly deliver events batched by session. One row per session is stored in a raw landing DLO, with each session’s events packed into a single JSON-array column. Downstream analysis use cases, such as funnel analysis, cohort segmentation, or identity-graph inputs, need one row per event.

Diagram showing one raw session row with a JSON-array events column expanding into multiple rows, one per event.

The solution: The parse-and-explode logic is worth pulling out of the entrypoint into a reusable helper. With Code Extensions, this can be done by putting files or modules in the payload/py-files/ directory, which needs to have this exact name.

1# payload/py-files/events.py
2from pyspark.sql.functions import col, explode_outer, from_json, to_timestamp
3from pyspark.sql.types import (
4    ArrayType, StringType, StructField, StructType,
5)
6
7
8EVENT_SCHEMA = StructType([
9    StructField("event_id", StringType(), True),
10    StructField("event_type", StringType(), True),
11    StructField("ts", StringType(), True),
12    StructField("path", StringType(), True),
13    StructField("value", StringType(), True),
14])
15EVENT_ARRAY_SCHEMA = ArrayType(EVENT_SCHEMA, True)
16
17# One row per event, carrying its parent session + user.
18def parse_events(df, session_col, user_col, events_col):
19    field = next(f for f in df.schema.fields if f.name == events_col)
20    events = (
21        col(events_col)
22        if isinstance(field.dataType, ArrayType)
23        else from_json(col(events_col).cast("string"), EVENT_ARRAY_SCHEMA)
24    )
25    return df.select(
26        col(session_col),
27        col(user_col),
28        explode_outer(events).alias("evt"),
29    ).select(
30        col("evt.event_id").alias("event_id__c"),
31        col(session_col),
32        col(user_col),
33        col("evt.event_type").alias("event_type__c"),
34        to_timestamp(col("evt.ts")).alias("event_ts__c"),
35        col("evt.path").alias("path__c"),
36        col("evt.value").alias("value__c"),
37    )

The entrypoint is now thinner: read, explode, filter, write.

1# payload/entrypoint.py
2from pyspark.sql.functions import col
3
4from datacustomcode.client import Client
5from datacustomcode.io.writer.base import WriteMode
6from events import parse_events
7
8
9def main():
10    client = Client()
11
12    sessions = client.read_dlo("User_Sessions__dll").select(
13        "session_id__c", "user_id__c", "events__c"
14    )
15    exploded = parse_events(sessions, "session_id__c", "user_id__c", "events__c")
16
17    output = exploded.filter(col("event_id__c").isNotNull())
18
19    client.write_to_dlo("User_Events__dll", output, WriteMode.OVERWRITE)
20
21
22if __name__ == "__main__":
23    main()

The results written to User_Events__dll are as expected for the sample dataset:

Query results showing User_Events__dll records with the expected, parsed-and-exploded values in columns like event_type__c and path__c, corresponding to the correct values in user_id__c and session_id__c.

You can find this entire example, including the sample data used, in the SDK’s docs/examples/script/json_explode directory.

Example 3: Score leads with a bundled scikit-learn model

The use case: The data scientist on your team trained a small scikit-learn classifier over categorical lead features (for example, industry, employee band, region, and source) and handed you the serialized artifact. You want a score on every lead as part of the nightly pipeline.

The solution: Similar to the py-files directory from the last example, you can ship other file types in the files directory. You can also ship external dependencies in Code Extensions by using the requirements.txt file, which in this example brings in scikit-learn, joblib, and pandas. Bundle the model as payload/files/lead_scorer.zip (a zip around lead_scorer.joblib — the joblib carries the sklearn pipeline plus feature schema and value domain). The training script is a plain scikit-learn Pipeline fit, run locally before deploying the Code Extension for any model adjustments.

The PySpark entrypoint can joblib.load the model, enumerate every combination of its categorical inputs, score that small grid with predict_proba, and left-join the grid against the leads DLO.

1import itertools
2import tempfile
3import zipfile
4from pathlib import Path
5
6import joblib
7import pandas as pd
8from pyspark.sql import Row
9from pyspark.sql.functions import coalesce, col, lit
10
11from datacustomcode.client import Client
12from datacustomcode.io.writer.base import WriteMode
13
14
15UNKNOWN = "__unknown__"
16MODEL_ARCHIVE = "lead_scorer.zip"
17MODEL_MEMBER = "lead_scorer.joblib"
18
19OUTPUT_COLUMNS = [
20    "id__c", "first_name__c", "last_name__c", "industry__c",
21    "employee_band__c", "region__c", "source__c", "score__c",
22]
23
24
25def load_bundled_model(client):
26    archive_path = client.find_file_path(MODEL_ARCHIVE)
27    extract_dir = Path(tempfile.mkdtemp(prefix="lead_scorer_"))
28    with zipfile.ZipFile(archive_path) as zf:
29        zf.extract(MODEL_MEMBER, path=extract_dir)
30    return joblib.load(extract_dir / MODEL_MEMBER)
31
32
33def build_score_grid(spark, bundle):
34    pipeline = bundle["pipeline"]
35    feature_cols = bundle["feature_cols"]
36    domain = bundle["feature_domain"]
37
38    combos = list(itertools.product(*(domain[c] for c in feature_cols)))
39    grid_pd = pd.DataFrame(combos, columns=feature_cols)
40    grid_pd["score__c"] = pipeline.predict_proba(grid_pd)[:, 1]
41
42    rows = [
43        Row(**{c: r[c] for c in feature_cols}, score__c=float(r["score__c"]))
44        for r in grid_pd.to_dict(orient="records")
45    ]
46    return spark.createDataFrame(rows), feature_cols
47
48
49def main():
50    client = Client()
51
52    leads = client.read_dlo("Lead__dll")
53    spark = leads.sparkSession
54
55    bundle = load_bundled_model(client)
56    grid, feature_cols = build_score_grid(spark, bundle)
57
58    normalized = leads
59    for c in feature_cols:
60        normalized = normalized.withColumn(c, coalesce(col(c), lit(UNKNOWN)))
61
62    scored = (
63        normalized.alias("l")
64        .join(grid.alias("g"), feature_cols, "left")
65        .withColumn("score__c", coalesce(col("score__c"), lit(0.0)))
66    )
67
68    client.write_to_dlo(
69        "Lead_Scored__dll", scored.select(*OUTPUT_COLUMNS), WriteMode.OVERWRITE
70    )
71
72
73if __name__ == "__main__":
74    main()

The grid is small, just a few thousand rows, so Spark’s planner auto-broadcasts it and every lead gets its score from a plain equi-join. This approach works well when features are naturally categorical or already banded (industry, region, tier, or segment). Same pattern for churn, propensity, and product-affinity: any model whose inputs are categorical or bucketable. coalesce to __unknown__ keeps the join total as new industries or sources appear.

The results written to Lead_Scored__dll are as expected for the sample dataset:

Query results showing Lead_Scored__dll records with score__c values calculated on each lead record.

Because Lead_Scored__dll is a first-class Data 360 object, once mapped into a DMO, an Agentforce agent like Coworker can query it directly: “Who are my top 10 hottest leads this week?”

You can find this entire example, including the sample data used, in the SDK’s docs/examples/script/lead_scoring directory.

Limits and other considerations

Code Extension jobs run under managed Spark with fixed memory, timeout, and dependency-size limits documented in the Considerations When Writing Code Extensions guide. Because local runs sample only the first 1000 rows (by default) via the Connect API, data-scale problems, such as partition skew on a deep account tree or an out-of-memory issue on a large explode, only surface on the deployed run, so test at full scale after deploy. The documented limitations will evolve over time based on the roadmap, which you can influence by reaching out to your account executive and discussing feature requests.

Conclusion

Code Extension is the right tool for the job when Python is the right language: imperative control flow, an existing library you want to reuse, a serialized model, a helper module worth unit-testing on its own. Code Extension for batch data transform is GA. What use cases does this unblock for you and your team? 

More Resources

About the author

Mark DeLaVergne is a Principal Software Engineer at Salesforce working on Code Extension in Data 360.  With 19 years in software engineering and architecture, he currently leads Code Extension and partners on wider Data 360 design.  You can find him on LinkedIn.

More Blog Posts

Run Complex Data Transformations in Data 360 with Code Extension

Run Complex Data Transformations in Data 360 with Code Extension

Build, deploy, run, and troubleshoot complex transformations with Python and PySpark while keeping execution governed by Data 360.September 10, 2026

The Salesforce Developer’s Guide to the Summer ’26 Release

The Salesforce Developer’s Guide to the Summer ’26 Release

Summer ’26 developer highlights: Hosted MCP Servers, LWC State Managers, Apex user-mode defaults, Agentforce Mobile SDK, and CLI updates with code examples.June 08, 2026

The Salesforce Developer’s Guide to the Winter ’27 Release

The Salesforce Developer’s Guide to the Winter ’27 Release

Winter ’27 developer highlights: the Salesforce Development plugin for Claude Code, Headless Experience Layer widgets, Angular and microfrontends in Multi-Framework, agents as MCP tools, LWC template expressions, Agent Script updates, and Web Console, with code examples.October 05, 2026