A Transformer pipeline using SQLTransformer objects (or any other Transformer) is a Spark solution which makes chaining transformations easy. For example:
from pyspark.ml.feature import SQLTransformer
from pyspark.ml import Pipeline, PipelineModel
df = spark.createDataFrame([
(0, 1.0, 3.0),
(2, 2.0, 5.0)
], ["id", "v1", "v2"])
sqlTrans = SQLTransformer(
statement="SELECT *, (v1 + v2) AS v3, (v1 * v2) AS v4 FROM __THIS__")
sqlSelectExpr = SQLTransformer(statement="SELECT *, (id * 2) AS v5 FROM __THIS__")
pipeline = Pipeline(stages=[sqlTrans, sqlSelectExpr])
pipelineModel = pipeline.fit(df)
pipelineModel.transform(df).show()
Another approach to chaining when all the transformations are simple expressions such as above, is to use a single SQLTransformer and string manipulations:
transforms = ['(v1 + v2) AS v3',
'(v1 * v2) AS v4',
'(id * 2) AS v5',
]
selectExpr = "SELECT *, {} FROM __THIS__".format(",".join(transforms))
sqlSelectExpr = SQLTransformer(statement=selectExpr)
sqlSelectExpr.transform(df).show()
Keep in mind that Spark SQL transformations can be optimized and will be faster than transforms defined as a Python User Defined Function (UDF).