Distributed Spark and Trino/Starburst adapter with SQL/JSON DAG execution and Spark writes to Trino, Iceberg and Hive.
Project description
trino-spark-adapter
trino-spark-adapter is a Python package for building hybrid Spark and Trino/Starburst data pipelines.
It provides:
- distributed Trino/Starburst reads executed from Spark executors;
- Spark writes to Trino through distributed
INSERTbatches; - Spark writes to Iceberg tables with Iceberg partition transforms and table properties;
- Spark writes to Hive-compatible tables or paths with Hive-style partitioning and bucketing;
- a DAG runner that executes
.sqland.jsonfiles by step number, running same-step files in parallel; - AES-CBC helpers backed by the
cryptographylibrary and Spark SQL UDF registration; - class-based logging with one logger per component.
Installation
pip install trino-spark-adapter==1.1.6
Core concepts
A DAG is a folder of dot-separated files executed by DagRunner. The naming pattern is:
{step}.{name}.{operation}.{extension}
| Segment | Role |
|---|---|
step |
Execution order. Files with the same step run in parallel. |
name |
Logical task name. Default Spark view for .trino_to_spark.json. |
operation |
Engine and bridge (spark, trino_to_spark, ...). |
extension |
File type (sql, json, ...). |
Example layout:
1.source.trino_to_spark.json # step 1, default view "source"
2.enrich.spark.sql # step 2
2.lookup.trino_to_spark.json # step 2, parallel with enrich, view "lookup"
2.a.trino.sql # step 2, short alias also works
2.b.spark.sql # step 2
3.export.spark_to_trino.json # step 3
For .trino_to_spark.json, the default Spark temporary view is the name segment unless target_view or view_name is set in the JSON. Example: 2.lookup.trino_to_spark.json creates view lookup.
| Suffix | Action |
|---|---|
.trino.sql |
Execute SQL statements on Trino/Starburst. |
.spark.sql |
Execute SQL statements on Spark. |
.trino_to_spark.json |
Load a Trino table or query into a Spark temporary view. |
.spark_to_trino.json |
Write a Spark view to Trino through distributed inserts. |
.spark_to_iceberg.json |
Write a Spark view to an Iceberg table with Spark Writer V2. |
.spark_to_hive.json |
Write a Spark view to a Hive-compatible table or path. |
Minimal DAG runner
from trino_spark_adapter import DagRunner, SparkAESHelper, TrinoConnectionConfig
spark = SparkAESHelper.get_spark(
app_name="trino_spark_adapter_job",
register_aes=True,
fail_if_missing_aes=False,
)
trino_config = TrinoConnectionConfig.from_env()
runner = DagRunner(
spark=spark,
trino_config=trino_config,
params={"calculation_date": "2026-01-05"},
reader_defaults={"fetch_size": 100_000, "num_ranges": 20, "num_partitions": 20},
)
results = runner.run_folder("dag")
Trino configuration
TrinoConnectionConfig.from_env() reads:
export TRINO_HOST="starburst.example.com"
export TRINO_USER="user"
export TRINO_PASSWORD="password"
export TRINO_ROLES="role"
export TRINO_VERIFY="false"
export TRINO_HTTP_SCHEME="https"
export TRINO_PORT="443"
Spark and AES helper
SparkAESHelper creates a Spark session from environment variables whose names start with spark. and optionally registers AES UDFs. Encryption runs in-process through the cryptography library (Cipher, algorithms.AES, modes.CBC). The key and IV are bound once at registration time and stored as private attributes on the UDF closure, so SQL only receives the column value (aes_decrypt(col)). Secrets stay in process memory and are not written to disk by the library.
export spark.app.name="trino_spark_adapter_job"
export spark.sql.shuffle.partitions="200"
export aes_key_str="...base64..."
export aes_iv_str="...base64..."
from trino_spark_adapter import SparkAESHelper
spark = SparkAESHelper.get_spark(register_aes=True)
spark.sql("SELECT aes_encrypt('abc') AS encrypted_value").show()
Registered functions:
aes_encrypt(value)aes_decrypt(value)
Date parameters
DateUtils can be used to compute placeholders used in DAG files.
from trino_spark_adapter import DateUtils
weekday_dates = DateUtils.generate_weekday_dates_between_start_stop(
start_dt="2026-01-01",
stop_dt="2026-02-01",
weekday=0,
)
du = DateUtils(weekday_dates[0])
params = du.to_params()
params.update({"calculation_date": du.today_tiret[:10]})
Typical generated keys include today_tiret, today_slash, last_day_tiret, last_week_tiret, last_month_tiret, last_quarter_tiret, last_semester_tiret, and last_year_tiret.
.trino_to_spark.json
Load a complete table:
{
"table_fullname": "catalog.schema.source_table",
"target_view": "source_table"
}
Distributed reads split on any column (DATE, TIMESTAMP, integers, decimals), not only partition columns. Provide either num_ranges or step_ranges, not both.
Equal number of date ranges:
{
"table_fullname": "catalog.schema.source_table",
"colname": "event_date",
"coltype": "DATE",
"format": "%Y-%m-%d",
"colname_start_value": "{start_date}",
"colname_stop_value": "{end_date}",
"num_ranges": 20
}
Fixed interval per range (pandas offset alias):
{
"table_fullname": "catalog.schema.source_table",
"colname": "event_ts",
"coltype": "TIMESTAMP",
"format": "%Y-%m-%d %H:%M:%S",
"colname_start_value": "2026-01-01 00:00:00",
"colname_stop_value": "2026-02-01 00:00:00",
"step_ranges": "7D"
}
Optional rounding (D, H, min, S) adjusts split boundaries. Numeric columns accept num_ranges or numeric step_ranges. Omit both to run one query on the full interval.
num_partitions controls Spark parallelism. num_ranges controls how many Trino queries are generated.
The runner creates or replaces the Spark temporary view named by target_view. When omitted, the view name is the second dot-separated segment of the file name. Example: 1.source.trino_to_spark.json creates view source.
.spark_to_trino.json
Write a Spark view to a Trino table. The target table is created automatically when it does not exist, based on the Spark schema.
{
"source_view": "prepared_view",
"target_table": "catalog.schema.target_table",
"repartition_by": ["entity_id"],
"num_partitions": 40,
"sort_by": ["entity_id"]
}
.spark_to_iceberg.json
Write a Spark view to an Iceberg table through the configured Spark Iceberg catalog.
{
"source_view": "prepared_view",
"catalog": "iceberg_catalog",
"schema": "analytics",
"table": "target_table",
"mode": "append",
"partition_spec": [
{"transform": "day", "column": "event_ts"},
{"transform": "bucket", "column": "entity_id", "num_buckets": 32}
],
"distribution_mode": "hash",
"format_version": 2,
"file_format": "PARQUET",
"repartition_by": ["entity_id"],
"num_partitions": 200
}
Common modes are create, replace, append, and overwrite_partitions.
.spark_to_hive.json
Write a Spark view to a Hive table:
{
"source_view": "prepared_view",
"table": "analytics.target_table",
"mode": "overwrite",
"format": "parquet",
"partition_by": ["event_date"],
"repartition_by": ["event_date"],
"num_partitions": 40
}
Write a Spark view to a path such as S3A:
{
"source_view": "prepared_view",
"path": "s3a://bucket/path/target_table",
"mode": "overwrite",
"format": "parquet",
"partition_by": ["event_date"]
}
Hive bucketing is supported only with table / saveAsTable:
{
"source_view": "prepared_view",
"table": "analytics.bucketed_table",
"mode": "overwrite",
"format": "parquet",
"bucket_by": ["entity_id"],
"num_buckets": 32,
"sort_by": ["entity_id"]
}
Logging
Every main component inherits from LogBase and exposes a class logger.
import logging
from trino_spark_adapter import DagRunner, DistributedTrinoSparkReader, SparkHiveWriter
DagRunner.logger().setLevel(logging.INFO)
DistributedTrinoSparkReader.logger().setLevel(logging.DEBUG)
SparkHiveWriter.logger().setLevel(logging.INFO)
The default formatter includes timestamp, class name and level:
[2026-01-05 09:15:12] [DagRunner] [INFO] START task type=trino_to_spark file=1.source.trino_to_spark.json
[2026-01-05 09:17:03] [DagRunner] [INFO] SUCCESS task type=trino_to_spark file=1.source.trino_to_spark.json elapsed=0:01:51
Debug-only expensive Spark actions should be guarded explicitly:
logger = DagRunner.logger()
if logger.isEnabledFor(logging.DEBUG):
logger.debug("row_count=%s", df.count())
Documentation
HTML documentation can be built locally:
pip install "trino-spark-adapter[docs]"
make -C docs html
Open docs/_build/html/index.html. The PyPI project page links to the online documentation URL.
Demos
Example notebooks are in demos/. They highlight DAG naming, parallel steps, default view names, date and numeric split strategies, writes, AES and date parameters.
Publishing
The source archive includes pypi.sh.
chmod +x .pypi.sh
./.pypi.sh
Resuming or partially executing a DAG
DagRunner.run_folder can execute only part of a DAG folder. This is useful
when a long run fails after several successful files and you want to restart
from the failing task.
runner.run_folder(
"dag",
start_from_file="6.transform.spark.sql",
)
You can also start from a zero-based index:
runner.run_folder(
"dag",
start_from_index=5,
)
For short validation runs, stop after a specific file or index:
runner.run_folder(
"dag",
stop_after_file="6.transform.spark.sql",
)
To run only a chosen subset, provide explicit file names. The files are still executed in the DAG folder's alphanumeric order:
runner.run_folder(
"dag",
files=[
"6.transform.spark.sql",
"7.export.spark_to_iceberg.json",
],
)
Optional Spark and Trino dependencies
Install core package:
pip install trino-spark-adapter==1.1.6
Optional extras:
pip install "trino-spark-adapter[spark]==1.1.4"
pip install "trino-spark-adapter[trino]==1.1.4"
pip install "trino-spark-adapter[all]==1.1.4"
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file trino_spark_adapter-1.1.6.tar.gz.
File metadata
- Download URL: trino_spark_adapter-1.1.6.tar.gz
- Upload date:
- Size: 48.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.4
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b0a10ab378f3269a90fe63be2759c0cbb142a033bb6602352b8159984195bf24
|
|
| MD5 |
ca4a65e995bb581caad1bed9b62f08e0
|
|
| BLAKE2b-256 |
04530b19f73f757191c579827a8f88bb1ef4823d1c42227c937b5b69e122ea66
|
File details
Details for the file trino_spark_adapter-1.1.6-py3-none-any.whl.
File metadata
- Download URL: trino_spark_adapter-1.1.6-py3-none-any.whl
- Upload date:
- Size: 44.5 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.4
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
68ff40c97700989adb8ac618f81e25b4ac958486c37e0a9b7b7bf0fb9a0bbd6a
|
|
| MD5 |
9ef52e719a62e2e977cc2723f9e96370
|
|
| BLAKE2b-256 |
62a79356355d5da2209180b3bc1cc4ac52d5e3e74493494ebec13ad075f7bebc
|