Skip to main content
Pre-release

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

dataforge-sdk

SDK for creating DataForge extensions.

Example projects and usage patterns: https://github.com/dataforgelabs/dataforge-sdk

Postgres Utilities

The dataforge.pg module provides helper functions to execute SQL operations against the DataForge Postgres metastore:

from dataforge.pg import select, update, pull

# Execute a SELECT query and return a Spark DataFrame
df = select("SELECT * FROM my_table")

# Execute an UPDATE/INSERT/DELETE query
update("UPDATE my_table SET col = 'value'")

# Trigger a new data pull for source_id 123
pull(123)

IngestionSession

The IngestionSession class manages a custom data ingestion process lifecycle.

from dataforge import IngestionSession

# Initialize a session (production use)
session = IngestionSession()

# Initialize a session (optional source_name/project_name for testing)
session = IngestionSession(source_name="my_source", project_name="my_project")

# Ingest data 
# pass a function returning a DataFrame (recommended to integrate logging with DataForge)
session.ingest(lambda: spark.read.csv("s3://bucket/path/input.csv"))

# pass a DataFrame (can be used for testing, not recommended for production deployment)
df = spark.read.csv("s3://bucket/path/input.csv")
session.ingest(df)

# Optional complete key snapshot for incremental custom ingestion.
# When omitted, df is treated as the complete current dataset for delete tracking.
# It must contain exactly the configured key columns and at least one row.
session.ingest(df, all_keys_df=all_current_keys_df)

# ingest empty dataframe to create 0-record input
session.ingest()


# Fail the process with error message
session.fail("Error message")

# Retrieve latest tracking fields
tracking = session.latest_tracking_fields()

# Retrieve connection parameters for the current source
connection_parameters = session.connection_parameters()

# Retrieve custom parameters for the current source
custom_parameters = session.custom_parameters()

# Retrieve system configuration for the current session
system_configuration = session.system_configuration

ParsingSession

The ParsingSession class manages a custom parse process lifecycle.

from dataforge import ParsingSession

# Initialize a session (production use)
session = ParsingSession()

# Initialize a session (optional input_id for testing)
session = ParsingSession(input_id=123)

# Retrieve custom parameters
params = session.custom_parameters()

# Retrieve system configuration for the current session
system_configuration = session.system_configuration

# Get the path of file to be parsed
path = session.file_path

# Run parsing: pass a DataFrame, a function returning a DataFrame or None (0-record file)
session.run(lambda: spark.read.json(session.file_path))

# Fail the process with error message
session.fail("Error message")

PostOutputSession

The PostOutputSession class manages a custom post-output process lifecycle.

from dataforge import PostOutputSession

# Initialize a session (production use)
session = PostOutputSession()

# Initialize a session (optional names for testing)
session = PostOutputSession(output_name="report", output_source_name="my_source", project_name="my_project")


# Get the path of file generated by preceding output process
path = session.file_path()

# Retrieve connection parameters for the current output
connection_parameters = session.connection_parameters()

# Retrieve custom parameters for the current source
custom_parameters = session.custom_parameters()

# Retrieve system configuration for the current session
system_configuration = session.system_configuration

# Run post-output logic: pass a function encapsulating custom code
session.run(lambda: print(f"Uploading file from {path}"))

# Fail the process with error message
session.fail("Error message")

Enqueue a separate custom ingestion

Configure the destination as an active batch custom ingestion source with its own notebook. In the orchestrated post-output notebook:

from dataforge import PostOutputSession

session = PostOutputSession()

def post_output():
    # Complete any writes needed by the destination before enqueueing.
    request = session.enqueue_ingestion("downstream_source")
    # {"status": "queued", "ingestion_queue_id": ..., "process_id": ..., "source_id": ...}

session.run(post_output)

The destination notebook runs independently with IngestionSession() and reads its own configured data. No nested ingestion session, project name, DataFrame, output view parameters, or snapshot is passed by enqueue_ingestion.

The source name is resolved in the post-output session's project (the same project as session.process.parameters['project_name']), without a fallback to another project. Missing, inactive, non-custom, non-batch, or unconfigured destinations raise an error. Call before the session closes, normally inside run's callback.

Every call immediately creates a queue record. Ingestion may start before post-output finishes, and later post-output failure or cancellation does not cancel the request. Repeated calls and notebook retries create additional queue records. Returning from this method means acceptance, not that ingestion has started or finished. Neither notebook waits for the other to finish. Complete the destination's data writes before calling enqueue_ingestion.

Queued requests wait until source, project, and environment initiation are enabled. They survive schedule reloads/restarts and never clear disable flags. If a destination is deactivated or loses its custom notebook configuration while waiting, it remains pending until that configuration is restored. Deleting the destination cancels its pending requests. Queue history records the originating process as post_output:<process_id>.

This uses the scheduler's refresh-latest behavior: pending requests for a destination can be served by one ingestion. It does not promise one ingestion per output event or preserve an output snapshot. Use durable, independently readable data and handle incremental progress in the destination notebook. Existing pg.pull retains its manual Pull Now behavior.

Release files for dataforge-sdk 11.0.0rc3

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dataforge-sdk 11.0.0rc3
File Size Uploaded
dataforge_sdk-11.0.0rc3.tar.gz 19.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dataforge-sdk 11.0.0rc3
File Interpreter ABI Platform
dataforge_sdk-11.0.0rc3-py3-none-any.whl Python 3 none any Details

Total release size: 42.1 kB

Release files / dataforge_sdk-11.0.0rc3.tar.gz

Download URL dataforge_sdk-11.0.0rc3.tar.gz
Size 19.0 kB
Tags Source
SHA-256 checksum
How to use checksums
055110987f94b32fc287985fa9b50eb9dc4f9c495354aefa37cccfbc06fb54e8
BLAKE2b-256 checksum
How to use checksums
c3e2027047a5643873af3ac2f5e21c5a1aeb9510c10f352cde0bc24568a182f9
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.12.14

Release files / dataforge_sdk-11.0.0rc3-py3-none-any.whl

Download URL dataforge_sdk-11.0.0rc3-py3-none-any.whl
Size 23.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
97b912e36946e85ffd7de98762ce96ccfe6193d4105611cd71a79e6dd06fb1e8
BLAKE2b-256 checksum
How to use checksums
f6892eb3252fa45d198d4189efa816ec7559c066b5d0e1a3591a15a8a99ba6da
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.12.14

Release history Release notifications | RSS feed

This release

11.0.0rc3 This release

2 release files

10.3.0

2 release files

10.2.0

2 release files

10.1.1

2 release files

10.1.0

2 release files

10.0.3

2 release files

10.0.2

2 release files

10.0.1

2 release files

9.2.6

2 release files

9.2.5

2 release files

9.2.4

2 release files

9.2.3

2 release files

9.2.2

2 release files

9.2.1

2 release files

9.2.0

2 release files

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