Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

SegmentStream pipeline SDK

segmentstream-pipeline is the small runtime contract between a workspace's Dagster definitions and the warehouse provisioned by SegmentStream. Dagster continues to own assets and jobs, while Ibis continues to own relational expressions. This package supplies lazy runtime configuration and durable warehouse I/O.

The initial connector supports BigQuery, automatic dataset creation, full-table replacement for unpartitioned assets, and native daily DATE partitioning. It reads the following non-secret configuration when a pipeline first accesses the warehouse:

  • SEGMENTSTREAM_WAREHOUSE_ENGINE
  • SEGMENTSTREAM_WAREHOUSE_CATALOG
  • SEGMENTSTREAM_WAREHOUSE_DEFAULT_NAMESPACE
  • SEGMENTSTREAM_WAREHOUSE_LOCATION (optional)

Configuration and authentication are deliberately lazy. Importing and validating definitions.py during a deployment build does not connect to a warehouse. In Cloud Run, the BigQuery connector uses the attached workload identity through Application Default Credentials.

import ibis
import ibis.expr.types as ir
import segmentstream.dagster as dg

from segmentstream import WAREHOUSE_IO_MANAGER_KEY, warehouse_resources


BRONZE_ORDERS = dg.AssetKey(["bronze", "orders"])
SILVER_ORDERS = dg.AssetKey(["silver", "orders"])


@dg.asset(
    key=BRONZE_ORDERS,
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
    kinds={"ibis"},
)
def orders() -> ir.Table:
    return ibis.memtable(
        [{"order_id": "o-1", "amount": 100.0}],
        schema={"order_id": "string", "amount": "float64"},
    )


@dg.asset(
    key=SILVER_ORDERS,
    ins={"orders": dg.AssetIn(key=BRONZE_ORDERS)},
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
    kinds={"ibis"},
)
def normalized_orders(orders: ir.Table) -> ir.Table:
    return orders.filter(orders.amount > 0)


defs = dg.Definitions(
    assets=[orders, normalized_orders],
    resources=warehouse_resources(),
)

Daily assets use Dagster's native daily partitions and declare the physical BigQuery DATE column through SegmentStream metadata:

from datetime import date

import ibis
import ibis.expr.types as ir
import segmentstream.dagster as dg

from segmentstream import (
    WAREHOUSE_IO_MANAGER_KEY,
    warehouse_asset_metadata,
)


daily = dg.DailyPartitionsDefinition(start_date="2026-01-01")


@dg.asset(
    key=["silver", "daily_orders"],
    partitions_def=daily,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
    metadata=warehouse_asset_metadata(partition_by_date="event_date"),
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
)
def daily_orders(context: dg.AssetExecutionContext) -> ir.Table:
    partition_dates = [date.fromisoformat(key) for key in context.partition_keys]
    return ibis.memtable(
        [{"event_date": value, "order_count": 0} for value in partition_dates],
        schema={"event_date": "date", "order_count": "int64"},
    )

SegmentStream accepts only unpartitioned assets and default-midnight DailyPartitionsDefinition assets with YYYY-MM-DD keys. Deployment inspection rejects other partition definitions and daily assets without physical partition metadata. The IO manager maps Dagster's partition time window to a half-open warehouse date range, filters upstream Ibis relations to that range, creates the table with native daily partitioning on first materialization, and atomically replaces only those dates on subsequent materializations.

Asset keys map to relations using a small convention:

  • ["orders"] uses the configured default namespace.
  • ["bronze", "orders"] uses the explicit bronze dataset.
  • Other key shapes are rejected.

The workspace project is always supplied by SegmentStream and cannot be overridden by an asset. Before writing an asset, the IO manager creates its validated dataset with CREATE SCHEMA IF NOT EXISTS in the configured location. This lets pipeline authors organize one workspace project into datasets such as bronze, silver, and gold without provisioning them separately.

For local package development, install this project in editable mode rather than adding a relative path dependency to a deployable pipeline.

Workspace pipelines declare only the SegmentStream SDK. It installs the pinned Dagster and Ibis versions that belong to that SDK release:

[project]
dependencies = [
  "segmentstream-pipeline[bigquery]==0.1.0a4",
]

Pipeline definitions import segmentstream.dagster as their curated Dagster namespace. Its objects are direct re-exports from Dagster, not wrappers. APIs outside that namespace are not part of the SegmentStream Cloud compatibility contract even if they remain importable from the underlying dependency.

Releases

Releases use the version declared in pyproject.toml and are published from the pipeline-sdk-v<version> Git tag by the protected pipeline-sdk-release.yml workflow. The workflow builds the wheel and source distribution in a job without publishing credentials, then uses PyPI Trusted Publishing from the pypi GitHub environment. No long-lived PyPI token is stored in GitHub.

PyPI releases are immutable. Increment the package version before creating a new release tag; do not reuse a version that has already been uploaded.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

segmentstream_pipeline-0.1.0a4.tar.gz (82.2 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

segmentstream_pipeline-0.1.0a4-py3-none-any.whl (17.9 kB view details)

Uploaded Python 3

File details

Details for the file segmentstream_pipeline-0.1.0a4.tar.gz.

File metadata

  • Download URL: segmentstream_pipeline-0.1.0a4.tar.gz
  • Upload date:
  • Size: 82.2 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for segmentstream_pipeline-0.1.0a4.tar.gz
Algorithm Hash digest
SHA256 ecba9ab3c02633357853aab2db4ef2fb8740131d09e49a813fdd8eb049f6d155
MD5 94311495d89e760c33c381013a9c08c0
BLAKE2b-256 e4df8ae869f4cea15c5182a865d9cc701b3fb341f2fb03e46c137bd647839ce2

See more details on using hashes here.

Provenance

The following attestation bundles were made for segmentstream_pipeline-0.1.0a4.tar.gz:

Publisher: pipeline-sdk-release.yml on segmentstream/segmentstream

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file segmentstream_pipeline-0.1.0a4-py3-none-any.whl.

File metadata

File hashes

Hashes for segmentstream_pipeline-0.1.0a4-py3-none-any.whl
Algorithm Hash digest
SHA256 6ff678c0df96419516b7fdfb988dd369521628a6a54b32f00e379f9cc078cd0f
MD5 76c8a0907c7fa9f28411f6bb049f8992
BLAKE2b-256 900995704a1a4f6fd4b46e80b5e802dfdb39517bf3f00656755f2abe1c41e0b2

See more details on using hashes here.

Provenance

The following attestation bundles were made for segmentstream_pipeline-0.1.0a4-py3-none-any.whl:

Publisher: pipeline-sdk-release.yml on segmentstream/segmentstream

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.
Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page