Skip to content
vim89Public

About

Let's be honest - most data pipeline frameworks treat types as suggestions. Config files are strings. Schemas are "validated" at runtime. Data quality is an afterthought. So, let's do differently

Topics

Resources

Contributing

Security policy

Stars

3 stars

Watchers

0 watching

Forks

Latest commit

 

History

245 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

flowforge - Type‑safe-first Data Engineering

Build

codecov Core Contracts Connectors Infrastructure

Release

Changelog Docs

Scala sbt JDK License

Build pipelines that won’t even compile when contracts drift. Keep transformations pure, put effects at the edges, and run on Spark and Flink.

Why (beliefs)

  • Runtime schema drift burns weekends. We believe failures should move left - into the compiler.
  • Side‑effects inside transforms amplify retries/speculation. We believe effects belong at the edges and must be idempotent.
  • Engineers deserve fast, local feedback. We believe pure transformations and compile‑fail tests make data engineering joyful again.

A story:

"A partner team removed a nullable column late Friday. We couldn’t roll back in time; both teams were up all night. If that change had been a compile error, we would have slept."

For Python/ETL folks (dbt/Airflow/Informatica/Talend):

Think "contracts like Pydantic/Avro - but enforced before jobs run," "pure functions you can test without a cluster," and "connectors/engines that make IO explicit and safe."

For EMs / Staff Data Architects:

You get compile‑time guarantees (not CI or runtime heuristics), a small opinionated surface, and batteries‑included defaults with escape hatches.

How (principles)

  • Compile‑time contracts: SchemaConforms[Out, Contract, Policy] proves compatibility; policies include Exact, Backward, Forward (+ Ordered/CI/ByPosition). See docs/how-it-fails.md. The engine behind it is ctdc-core, a standalone library; com.flowforge.core.contracts aliases its types so there is one implementation to maintain.
  • Typestate builder: build() exists only when source, transforms, and sink are present. Incomplete pipelines are unbuildable.
  • Pure vs effect boundary: transforms are pure functions; F[_] only at IO edges; engines plug into a single algebra.
  • Pictures over prose: see flowchart.svg and optionality.md.

What (The framework)

  • Core: contracts, builder, EffectSystem, DataAlgebra.
  • Engines: Spark. Flink is pinned to Scala 2.12 and does not build today, see the Flink note below.
  • Connectors: filesystem, GCS, and JDBC through Spark's own JDBC source. See docs/connectors/CAPABILITIES.md for what each one supports.
  • Data Quality: native checks by default; optional Deequ when present.
  • Template: flowforge.g8 for new projects.

Diagrams (pictures > words)

Compile‑time contracts flow

Field vs Element Optionality

Scala 2 Magnolia UML

Scala 3 Mirrors UML

Quick links

Module status (coverage)

  • Core: Core Coverage
  • Contracts: Contracts Coverage
  • Connectors: Connectors Coverage
  • Infrastructure: Infrastructure Coverage

Guarantees (Non‑negotiables)

  • Compile‑fail contracts for typed endpoints under policy lattice
  • Typestate builder: build() only when complete - incomplete pipelines can’t compile
  • Pure transforms; effectful edges; idempotent side‑effects by design
  • See: docs/design/framework-behaviors.md

10‑Minute quickstart

Prereq: JDK 17+, sbt 1.9+

1) Clone & build

git clone https://github.com/vim89/flowforge.git && cd flowforge
sbt compile

2) See a compile‑time contract failure (red → green)

// Paste in a scratch file under modules/examples and run `sbt examples/compile`
import com.flowforge.core.contracts._
final case class Out(id: Long)
final case class Contract(id: Long, email: String)
implicitly[SchemaConforms[Out, Contract, SchemaPolicy.Exact]] // compile error: missing email

Green needs the shape to satisfy the contract. Backward allows a producer to carry extra fields and to omit a field the contract declares optional, so it does not rescue a missing required field:

final case class OutWithExtra(id: Long, email: String, tag: String)
implicitly[SchemaConforms[OutWithExtra, Contract, SchemaPolicy.Backward]] // ok: tag is extra

final case class OptionalEmail(id: Long, email: Option[String])
implicitly[SchemaConforms[Out, OptionalEmail, SchemaPolicy.Backward]] // ok: email is optional

3) Build a pipeline - typestate forbids incomplete builds

import cats.effect.IO
import com.flowforge.core.PipelineBuilder
import com.flowforge.core.contracts._
import com.flowforge.core.instances.EffectInstances._ // brings EffectSystem[IO]
import com.flowforge.core.types._

final case class User(id: Long, email: String)
val src  = TypedSource[User](LocalDataSource("/tmp/in", DataFormat.Parquet))
val sink = TypedSink[User](LocalDataSink("/tmp/out", DataFormat.Parquet))

PipelineBuilder[IO]("demo")
  .addTypedSource[User, User, SchemaPolicy.Exact](src, _ => IO.pure(User(1, "a@b")))
  .noTransform
  .addTypedSink[User, SchemaPolicy.Exact](sink, (_, _) => IO.unit)
  .build() // build is available only now

4) Explore diagrams and failure messages

Quickstart paths

Path Goal Commands
A - Examples Run a pipeline locally (no cluster) sbt "examples/runMain com.flowforge.examples.SimpleGoldenPath"
B - Red→Green See compile‑time error then fix Use the snippet above; run sbt compile
C - Spark path Exercise the Spark algebra sbt engines-spark/test (Delta and SCD tests are opt‑in integration tests)
D - New project Scaffold with g8 The template lives in this repo: sbt new file://$PWD/flowforge.g8

Compatibility

The versions CI runs are listed in docs/reference/compatibility.md. That file is the only place this repo states a tested version, so the list is not repeated here.

Flink (2.12)

Flink's Scala API is 2.12-only, so engines-flink is pinned to Scala 2.12. core and connectors publish 2.13 and 3 only, so the module cannot resolve its own dependencies, and no CI job builds it. Treat Flink as unfinished work rather than an engine you can pick today. The gap is tracked in docs/plan/v1.0-readiness.md.

Architecture (at a glance)

The diagrams above summarize derivation and policy checks; see also docs/diagrams/compile-time-contracts/guide.md for narrative.

Examples & demos

  • Examples module: modules/examples (runnable demos)
  • Optional Deequ mode: add -Dff.quality.mode=deequ (auto‑enables when on classpath)

Documentation map

Release & versioning

FAQ

  • Scala 3?
    • Core compiles on Scala 3; engines depend on Spark/Flink ecosystem (Spark 3.x limits Scala 3 today).
  • Why compile‑time vs tests?
    • Tests are sampled; compile‑time proofs are exhaustive for shapes and policy compatibility.
  • How does this compare to Databricks DLT/Dagster/dbt?
    • They perform runtime/CI checks; FlowForge enforces compile‑time gates and typestate builder. See docs/evidence for deeper comparisons.

Contributing

We welcome folks from Python/ETL backgrounds and JVM veterans alike. Start with docs/contributing/HANDBOOK.md, then pick an issue. Please run sbt scalafmtAll and sbt "scalafixAll" before submitting.

License

AGPLv3


Flowforge Hybrid Licensing Model

Flowforge adopts a hybrid licensing structure combining open innovation and IP protection.

  • Legacy releases stay under the license they shipped with. Tags v0.7.0, v0.8.0, v0.9.0-rc.1 and v0.9.0 are MIT; v0.8.1 is AGPLv3.
  • Active and future releases (v1.0 and onward) are licensed under AGPLv3. LICENSE is the unmodified AGPLv3 text and adds no further restrictions.
  • Commercial usage (offering as SaaS, embedding in proprietary systems, or internal closed-source deployments) requires a separate commercial license. See COMMERCIAL_LICENSE.md for template.
  • Contributor License Agreement (CLA) in CLA.md governs contribution terms, ensuring compatibility with the hybrid licensing framework.
  • Commercial exceptions and dual-licensing are handled directly by Vitthal Mirji for partners and enterprise use.

The goal: protect Flowforge’s compile-time innovation while keeping community use free and open.

About

Let's be honest - most data pipeline frameworks treat types as suggestions. Config files are strings. Schemas are "validated" at runtime. Data quality is an afterthought. So, let's do differently

Topics

Resources

Contributing

Security policy

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages