Versions Compared


  • This line was added.
  • This line was removed.
  • Formatting was changed.


This change addresses the use-cases described in the motivation section, where we need to compose an Estimator from a DAG of Estimator/Transformer. Note that the Graph/GraphBuilder supports Estimator class whose input schemas are different from its fitted Transformer.

4) Added transformSchemas()  to the Stage interface.

This is needed to validate the compatibility of input schemas with a given Estimator/Transformer instance.

5) Added setStateStreams and getStateStreams to the Transformer interface.

This change addresses the use-cases described in the motivation section, where a running Transformer needs to ingest the model state streams emitted by a Estimator, which could be running on a different machine.

65) Removed the methods PipelineStage::toJson and PipelineStage::loadJson. Add methods save(...) and load(...) to the Stage interface.

The following changes are relatively minor:

76) Removed TableEnvironment from the parameter list of fit/transform APIs.

This change simplifies the usage of fit/transform APIs.

87) Added PipelineModel and let Pipeline implement only the Estimator. Pipeline is no longer a Transformer.

This change makes the experience of using Pipeline consistent with the experience of using Estimator/Transformer, where a class is either an Estimator or a Transformer.

98) Removed Pipeline::appendStage from the Pipeline class.

This change makes the concept of Pipeline consistent with that of Graph/GraphBuilder. Neither Graph nor Pipeline provides the API to construct themselves.

109) Removed the Model interface.

This change simplifies the class hierarchy by removing a redundant class. It follows the philosophy of only adding complexity when we have explicit use-case for it.

1110) Renamed PipelineStage to Stage and add the PublicEvolving tag to the Stage interface.


Code Block
 * Base class for a stage in a Pipeline or Graph. The interface is only a concept, and does not have any actual
 * functionality. Its subclasses could be Estimator, Transformer or MLfunc. No other classes should inherit this
 * interface directly.
 * <p>Each stage is with parameters, and requires a public empty constructor for restoration.
 * @param <T> The class type of the Stage implementation itself.
 * @see WithParams
interface Stage<T extends Stage<T>> extends WithParams<T>, Serializable {
     * This method checks the compatibility between input schemas, stage's parameters and stage's
     * logic. It should raise an exception if there is any mismatch, e.g. the number of input
     * schemas is wrong, or if a required field is missing from a schema.
     * <p>If there is no mismatch, the method derives and returns the output schemas from the input
     * schemas.
     * <p>Note that the output schemas of a given Estimator instance should equal the output schemas
     * of the Transformer instance fitted by this Estimator instance, suppose the same list of input
     * schemas are used as inputs to the fit/transform methods respectively.
     * <p>For any Estimator instance added in a Pipeline, the transform method of the Transformer returned by this
     * Estimator instance should be able to accept a list of tables of the same length and schemas as the fit method of
     * this Estimator instance.
     * @param schemas the list of schemas of the input tables.
     * @return the list of schemas of the output tables.
    TableSchema[] transformSchemas(TableSchema... schemas);
 * @param <T> The class type of the Stage implementation itself.
 * @see WithParams
interface Stage<T extends Stage<T>> extends WithParams<T>, Serializable {
     * Saves this stage to the given path.
    void save(String path);

     * Loads this stage from the given path.
    void load(String path);

 * A MLFunc is a Stage that takes a list of tables as inputs and produces a list of
 * tables as results. It can be used to encode a generic multi-input multi-output machine learning function.
 * @param <T> The class type of the MLFunc implementation itself.
public interface MLFunc<T extends MLFunc<T>> extends Stage<T> {
     * Applies the MLFunc on the given input tables, and returns the result tables.
     * @param inputs a list of tables
     * @return a list of tables
    Table[] transform(Table... inputs);

 * A Transformer is a MLFunc with additional support for state streams, which could be set by the Estimator that fitted
 * this Transformer. Unlike MLFunc, a Transformer is typically associated with an Estimator.
 * @param <T> The class type of the Transformer implementation itself.
public interface Transformer<T extends Transformer<T>> extends MLFunc<T> {
     * Uses the given list of tables to update internal states. This can be useful for e.g. online
     * learning where an Estimator fits an infinite stream of training samples and streams the model
     * diff data to this Transformer.
     * <p>This method may be called at most once.
     * @param inputs a list of tables
    default void setStateStreams(Table... inputs) {
        throw new UnsupportedOperationException("this method is not implemented");

     * Gets a list of tables representing changes of internal states of this Transformer. These
     * tables might come from the Estimator that instantiated this Transformer.
     * @return a list of tables
    default Table[] getStateStreams() {
        throw new UnsupportedOperationException("this method is not implemented");

 * An Estimator is a Stage that takes a list of tables as inputs and produces a Transformer.
 * @param <E> class type of the Estimator implementation itself.
 * @param <M> class type of the Transformer this Estimator produces.
public interface Estimator<E extends Estimator<E, M>, M extends Transformer<M>> extends Stage<E> {
     * Trains on the given inputs and produces a Transformer.
     * @param inputs a list of tables
     * @return a Transformer
    M fit(Table... inputs);

 * A Pipeline acts as an Estimator. It consists of an ordered list of stages, each of which could be
 * an Estimator, Transformer or MLFunc.
public final class Pipeline implements Estimator<Pipeline, PipelineModel> {

    public Pipeline(List<Stage<?>> stages) {...}

    public PipelineModel fit(Table... inputs) {...}

    /** Skipped a few methods, including the implementations of the Estimator APIs. */

 * A PipelineModel acts as a Transformer. It consists of an ordered list of Transformers or MLFuncs.
public final class PipelineModel implements Transformer<PipelineModel> {

    public PipelineModel(List<Transformer<?>> transformers) {...}

    /** Skipped a few methods, including the implementations of the Transformer APIs. */


Code Block
 * A Graph acts as an Estimator. A Graph consists of a DAG of stages, each of which could be
 * an Estimator, Transformer or MLFunc. When `Graph::fit` is called, the stages are executed in a
 * topologically-sorted order. If a stage is an Estimator, its `Estimator::fit` method will be
 * called on the input tables (from the input edges) to fit a model. Then the model, which is a
 * Transformer, will be used to transform the input tables to produce output tables to the output
 * edges. If a stage is a Transformer or MLFunc, its `MLFunc::transform` method will be called on the
 * input tables to produce output tables to the output edges. The fitted model from a Graph is a
 * GraphModel, which consists of fitted models and transformers, corresponding to the Graph's
 * stages.
public final class Graph implements Estimator<Graph, GraphModel> {
    public Graph(...) {...}

    public GraphModel fit(Table... inputs) {...}

    public TableSchema[] transformSchemas(TableSchema... schemas) {
        return schemas;

    /** Skipped a few methods, including the implementations of some Estimator APIs. */

 * A GraphModel acts as a Transformer. A GraphModel consists of a DAG of Transformers or MLFuncs. When
 * `GraphModel::transform` is called, the stages are executed in a topologically-sorted order. When
 * a stage is executed, its `MLFunc::transform` method will be called on the input tables (from
 * the input edges) to produce output tables to the output edges.
public final class GraphModel implements Transformer<GraphModel> {
    /** Skipped a few methods, including the implementations of the Transformer APIs. */

 * A GraphBuilder provides APIs to build Graph and GraphModel from a DAG of Estimator, Transformer and MLFunc instances.
public final class GraphBuilder {
     * Specifies the upper bound (could be loose) of the number of output tables that can be
     * returned by the Transformer::getStateStreams and Transformer::transform methods, for any
     * stage involved in this Graph.
     * <p>The default upper bound is 20.
    public GraphBuilder setMaxOutputLength(int maxOutputLength) {...}

     * Creates a TableId associated with this GraphBuilder. It can be used to specify the passing of
     * tables between stages, as well as the input/output tables of the Graph/GraphModel generated
     * by this builder.
    public TableId createTableId() {...}

     * If the stage is an Estimator, both its fit method and the transform method of its fitted Transformer would be
     * invoked with the given inputs when the graph runs.
     * <p>If this stage is a Transformer or MLfunc, its transform method would be invoked with the given
     * inputs when the graph runs.
     * Returns a list of TableIds, which represents outputs of the Transformer::transform invocation.
    public TableId[] getOutputs(Stage<?> stage, TableId... inputs) {...}

     * If this stage is an Estimator, its fit method would be invoked with estimatorInputs, and the transform method
     * of its fitted Transformer would be invoked with transformerInputs, when the graph runs.
     * <p>This method throws Exception if the stage is a Transformer or MLFunc.
     * This method is useful when the state is an Estimator AND the Estimator::fit needs to take a different list of
     * Tables from the Transformer::transform of the fitted Transformer.
     * Returns a list of TableIds, which represents outputs of the Transformer::transform invocation.
    public TableId[] getOutputs(Stage<?> stage, TableId[] estimatorInputs, TableId[] transformerInputs) {...}

     * The GraphModel::setStateStreams should invoke the setStateStreams of the corresponding stage
     * with the corresponding inputs.
    void setStateStreams(Stage<?> stage, TableId... inputs) {...}

     * The GraphModel::getStateStreams should invoke the getStateStreams of the corresponding stage.
     * <p>Returns a list of TableIds, which represents outputs of the getStateStreams invocation.
    TableId[] getStateStreams(Stage<?> stage) {...}

     * Returns a Graph instance which the following API specification:
     * 1) Graph::fit should take inputs and returns a GraphModel with the following specification.
     * 2) GraphModel::transform should take inputs and return outputs.
     * 3) GraphModel::setStateStreams should take inputStates.
     * 4) GraphModel::getStateStreams should return outputStates.
     * The fit/transform/setStateStreams/getStateStreams should invoke the APIs of the internal stages in the order specified by the DAG of stages.
    Graph build(TableId[] inputs, TableId[] outputs, TableId[] inputStates, TableId[] outputStates) {...}

     * Returns a Graph instance which the following API specification:
     * 1) Graph::fit should take estimatorInputs and returns a GraphModel with the following specification.
     * 2) GraphModel::transform should take transformerInputs and return outputs.
     * 3) GraphModel::setStateStreams should take inputStates.
     * 4) GraphModel::getStateStreams should return outputStates.
     * The fit/transform/setStateStreams/getStateStreams should invoke the APIs of the internal stages in the order specified by the DAG of stages.
     * This method is useful when the Graph::fit needs to take a different list of Tables from the GraphModel::transform of the fitted GraphModel.
    Graph build(TableId[] estimatorInputs, TableId[] transformerInputs, TableId[] outputs, TableId[] inputStates, TableId[] outputStates) {...}

     * Returns a GraphModel instance which the following API specification:
     * 1) GraphModel::transform should take inputs and returns outputs.
     * 2) GraphModel::setStateStreams should take inputStates.
     * 3) GraphModel::getStateStreams should return outputStates.
     * The transform/setStateStreams/getStateStreams should invoke the APIs of the internal  stages in the order specified by the DAG of stages.
     * This method throws Exception if any stage of this graph is an Estimator.
    GraphModel buildModel(TableId[] inputs, TableId[] outputs, TableId[] inputStates, TableId[] outputStates) {...}

    // The TableId is necessary to pass the inputs/outputs of various API calls across the
    // Graph/GraphModel stags.
    static class TableId {}

