Run data pipelines in production with ease.
Design processing flows on a visual canvas and deploy them to managed environments. Built from experience with real-time data streaming. It supports a wide variety of sources and sinks, and the catalog keeps growing.
Compose pipelines on a visual canvas
Drag sources, transforms, conditionals and sinks onto a canvas and connect them. Branch on field values, reshape records in Java, and route to multiple destinations.
- Kafka, database and HTTP sources & sinks, transforms and conditional branches.
- Schema-based steps: record types validated as you build.
- Versioned designs you can review before deploying.
topic: orders
Enrich order
region == "EU"
postgres · eu
postgres · global
PIPELINE · order-enrichment
topic: orders
via Schema Registry
Transformation.java
region == "EU"
Postgres · orders_eu
discard · log skipped
A comprehensive step catalog for production pipelines
Consume from a Source, decode with the right deserializer (String, Avro with or without Schema Registry, JSON), reshape records with a custom Java Transformation, branch on field values with If, and land the result in a Sink, or Skip what you don't need. Each step's output schema automatically becomes the next step's input, so the whole flow stays type checked from source to sink.
Write Java without the project setup
TypedFlows generates the input schema and a typed record interface from the previous step. You write plain Java against typed records, in the browser, with compiler feedback as you type.
- Input schema & typed record interface generated for you.
- Live compilation with error markers.
- You define the output schema, and downstream steps see it instantly.
- No jars to build, no local toolchain, no cluster round-trips.
package com.typedflows.usercode;import com.typedflows.usercode.input.InputMain;import com.typedflows.usercode.output.Main;import java.time.LocalDateTime;import java.time.temporal.ChronoUnit;/// Class representing a set of methods and fields for transformation of a single incoming record into the/// defined output data structure./// This is the user's task to provide the full required transformation code.public class Transformation { /// Transforms one incoming record to one outgoing record. Main.$Factory has a method for the output /// record and for every nested record type in the output schema. Objects it returns start empty, so the /// user's task is to provide the required logic and calculations to achieve the desired output. There is no /// automatic mapping here, therefore all fields from the input, which are expected to remain unchanged in /// the output, must be properly mapped by the user within this operation. Use setters to fill them. /// Returning the structure unpopulated is what causes the mandatory field check to fail, which reports /// "Output data validation error. Missing fields: ..." and stops the record instead of passing it on. public Main transform(InputMain inputDataStructure, Main.$Factory factory) { Main outputDataStructure = factory.newMain(); InputMain.SupplySituation inputSupplySituation = inputDataStructure.getValue(); // Delivery date is given as a string String nextDeliveryString = inputSupplySituation.getRequirements().getNextDelivery(); LocalDateTime nextDeliveryDateTime = LocalDateTime.parse(nextDeliveryString); // To ensure test stability, instead of now(), hardcoded date time is used LocalDateTime pseudoNow = LocalDateTime.of(2025, 1, 1, 12, 0); int productionDaysToNextDelivery = Math.abs(Math.toIntExact(ChronoUnit.DAYS.between(nextDeliveryDateTime, pseudoNow))); int requiredMaterials = productionDaysToNextDelivery * inputSupplySituation.getRequirements().getRequirementPerDay(); int totalQuantity = calculateTotalQuantity(inputSupplySituation.getStorage()); int storageAtDeliveryDay = totalQuantity - requiredMaterials; Main.SupplyStatus supplyStatus = determineSupplyStatus(storageAtDeliveryDay); // To ensure test stability, instead of now(), hardcoded date time is used (1 minute later) LocalDateTime pseudoCalcTimestamp = LocalDateTime.of(2025, 1, 1, 12, 1); // Building output structure Main.SupplySummaryKey supplySummaryKey = factory.newSupplySummaryKey(); supplySummaryKey.setPlantId(inputDataStructure.getKey()); supplySummaryKey.setMaterialId(inputDataStructure.getValue().getMaterialId()); Main.SupplySummaryValue supplySummaryValue = factory.newSupplySummaryValue(); supplySummaryValue.setPlantId(inputDataStructure.getKey()); supplySummaryValue.setMaterialId(inputDataStructure.getValue().getMaterialId()); supplySummaryValue.setTotalStorage(totalQuantity); supplySummaryValue.setNextDelivery(nextDeliveryDateTime); supplySummaryValue.setCalculationTimestamp(pseudoCalcTimestamp); supplySummaryValue.setSupplyStatus(supplyStatus); supplySummaryValue.setStorageAtDeliveryDay(storageAtDeliveryDay); outputDataStructure.setKey(supplySummaryKey); outputDataStructure.setValue(supplySummaryValue); return outputDataStructure; } private int calculateTotalQuantity(InputMain.StorageData storageData) { int plantQuantity = storageData.getPlantQuantity(); int storageQuantity = storageData.getStorageQuantity(); return plantQuantity + storageQuantity; } private Main.SupplyStatus determineSupplyStatus(int storageAtDeliveryDay) { if (storageAtDeliveryDay > 0) { return Main.SupplyStatus.RESERVE; } if (storageAtDeliveryDay == 0) { return Main.SupplyStatus.AS_PLANNED; } return Main.SupplyStatus.SHORTAGE; }}Flexible database steps
Read from a table or write to one, through JDBC connections defined once per environment. A Database Source polls a table into your pipeline; a Database Sink lands records with a query you control.
- Connections managed per environment, referenced by name from the steps.
- Generate from DDL: introspect the table and get schema + code for free.
- Poll strategies for change capture without touching the source system.
insert("orders_eu").values(record).execute();Your Java builds the query. The connection comes from the environment.
A flow you might actually ship
Marketing wants customer changes from the CRM database as events. This flow reads the table, shapes each row into an event, keeps the customers who gave consent and publishes them to Kafka. Four steps, no glue code.
crm-db · public.customerspoll: updated_at + id
row → CustomerEventnormalise email, add region
marketing_consent== true
topic: customer-events
discard quietly
Promote from dev to prod with confidence
Each environment bundles its own Kafka clusters, schema registries and database connections. Build in development, validate in staging, then promote the exact same design. TypedFlows re-resolves every connection against the target environment, so nothing is rewritten by hand.
- Per-environment Kafka clusters & schema registries.
- Database connections managed centrally and reused by sinks.
- One design, promoted across environments.
Change one step at deploy time
A deployment runs a design in an environment. When a single step has to read or write somewhere else, override that step's configuration: its Kafka cluster, schema registry, topic or consumer group. Every other step keeps the environment's settings, and the design stays as it is.
- Tick the fields to override; the rest keep the design's default value.
- Point a source at another cluster or topic, for a replay or a migration.
- Run the same design twice with separate consumer groups.
- Overrides belong to the deployment, so the design stays reusable.
Everything a data pipeline needs
One platform for the whole path from source to sink.
Visual flow design
Build pipelines by connecting nodes.
Managed deployments
Promote a design and TypedFlows runs and supervises it.
Custom Java transformations
Write plain Java against a generated input schema, compiled live.
Environment promotion
Dev → staging → prod with isolated infrastructure.
Schema-based
Every step declares and validates its input schema.
Kafka & database connectors
Topics and databases as sources and sinks: wire format handled correctly.
What teams build with TypedFlows
If it moves between systems and needs to land somewhere reliable, it's a fit.
Enrich & reshape in flight
Read from topics, map and enrich records against a schema, and write the clean result downstream.
Topics to tables, and back
Land Kafka messages into Postgres and other stores, or poll a table into a topic, with conditional routing on the way.
Validation in the pipeline
Enforce schemas across environments so bad data never reaches a consumer.
Free toolkit for data engineers
Standalone browser tools from the team behind TypedFlows: validate, infer, encode and decode Avro & Kafka data.