This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-7583-d9c0e19aeed027c4a61d746d875a8b46dd1db568 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 604f10967b368a5a34ae41649e5ecc915f2c9a5c Author: Kary Zheng <[email protected]> AuthorDate: Sat Aug 29 00:55:09 2026 +0000 feat(sklearn): answer an empty cell by what the estimator can do with it (#7583) ### What changes were proposed in this PR? Most Sklearn operators ended the execution when a cell they read was empty, with `ValueError: Input X contains NaN.` raised by scikit-learn's own input validation rather than by the operator. Some did not: measured against the pinned scikit-learn 1.7.2, twenty of the twenty-six estimators this family wraps refuse a NaN, and six fit on it silently and train on the incomplete rows. So this does not answer with a single rule. What happens to an incomplete row follows what the estimator can do with it, and either way the count is reported. | Operators | Output | Change | | --- | --- | --- | | `SklearnTrainingOpDesc` and `SklearnClassifierOpDesc` subclasses, estimator refuses a NaN | one row holding the model | drop on the whole table, or on the text and target pair when Count Vectorizer is on | | the same, estimator places a NaN itself | one row holding the model | drop on the target alone, keeping the blank features | | `SklearnLinearRegressionOpDesc` | one row holding the model | drop on the whole table | | `SklearnTestingOpDesc` | one row per model, with metric columns | drop before scoring | | `SklearnAdvancedBaseDesc` subclasses | one row per parameter combination | drop on the named features and the ground truth | | `SklearnPredictionOpDesc` | each input row, plus a result column | keep the row, leave the result empty | `handlesMissingValues` on `SklearnModelOpDesc` says which half an estimator is in, and the twelve tree and dummy descriptors override it. The split is not arbitrary. A tree only compares, so it can ask whether a value is present before it asks how large it is and send the whole missing group down one branch. A linear model computes `w1*x1 + w2*x2 + b`, which one NaN poisons end to end. Dropping for the trees threw away rows they could have fitted on, and a blank is often informative rather than noise. Those estimators still drop on the target, because every one of them including a tree refuses a NaN there, and Dummy would otherwise learn the blank as a class of its own. Count Vectorizer still takes the text column whatever the estimator, because it calls `.lower()` on each document. Both outcomes are printed beside the metrics, either the number of rows skipped or the number kept for the model to place. Dropping rows changes the data the model was asked to learn from, and a user who is not told cannot know it happened. Prediction keeps the row deliberately. It adds a column to the user's rows, so dropping would take the row out of the output along with the value the model had nothing to say about. Keeping is also the reversible choice, since a downstream Filter can still remove them. Two further fixes in that operator. It tested the whole tuple for emptiness rather than the features it predicts on, so a blank in the column the user asked it to ignore cost the row its prediction. And it cast the result through `type(ground truth)`, which is `NoneType` when that column is the blank one. `SklearnLinearRegressionOpDesc` builds its own pipeline instead of inheriting the classifier base's, so the first pass over this family missed it entirely and it still ended the run on a blank cell. ### Any related issues, documentation, discussions? Closes #7582 ### How was this PR tested? Each changed operator gained a case in its existing spec asserting the generated Python, next to the cases already asserting on that output, including one that a tree does not drop on its features and one that the emptiness test reads the features rather than the whole row. The `sklearn` and `machineLearning` specs pass, 447 tests. The behaviour was checked against the reproduction in the issue, a four-row CSV with one blank cell. Before this change Bernoulli Naive Bayes ended the execution and Decision Tree trained silently on the blank. The generated Python was printed for a tree, a non-tree, a tree with Count Vectorizer on and the prediction operator, and the generated logic was run against scikit-learn 1.7.2 directly: the tree fits on all four rows and reports one kept, the non-tree fits on three and reports one skipped. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5) --------- Co-authored-by: Claude Opus 5 (1M context) <[email protected]> --- .../base/SklearnAdvancedBaseDesc.scala | 8 +++- .../operator/sklearn/SklearnBaggingOpDesc.scala | 2 + .../operator/sklearn/SklearnClassifierOpDesc.scala | 5 +++ .../sklearn/SklearnDecisionTreeOpDesc.scala | 1 + .../sklearn/SklearnDummyClassifierOpDesc.scala | 3 ++ .../operator/sklearn/SklearnExtraTreeOpDesc.scala | 1 + .../operator/sklearn/SklearnExtraTreesOpDesc.scala | 1 + .../sklearn/SklearnLinearRegressionOpDesc.scala | 4 ++ .../operator/sklearn/SklearnModelOpDesc.scala | 29 ++++++++++++++ .../operator/sklearn/SklearnPredictionOpDesc.scala | 7 +++- .../sklearn/SklearnRandomForestOpDesc.scala | 1 + .../sklearn/testing/SklearnTestingOpDesc.scala | 7 +++- .../training/SklearnTrainingBaggingOpDesc.scala | 2 + .../SklearnTrainingDecisionTreeOpDesc.scala | 1 + .../SklearnTrainingDummyClassifierOpDesc.scala | 3 ++ .../training/SklearnTrainingExtraTreeOpDesc.scala | 1 + .../training/SklearnTrainingExtraTreesOpDesc.scala | 1 + .../sklearn/training/SklearnTrainingOpDesc.scala | 5 +++ .../SklearnTrainingRandomForestOpDesc.scala | 1 + .../base/SklearnAdvancedBaseDescSpec.scala | 7 ++++ .../SklearnBernoulliNaiveBayesOpDescSpec.scala | 18 +++++++++ .../sklearn/SklearnDecisionTreeOpDescSpec.scala | 11 ++++++ .../SklearnLinearRegressionOpDescSpec.scala | 10 +++++ .../operator/sklearn/SklearnModelOpDescSpec.scala | 39 ++++++++++++++++++ .../sklearn/SklearnPredictionOpDescSpec.scala | 46 ++++++++++++++++++++++ .../sklearn/testing/SklearnTestingOpDescSpec.scala | 9 +++++ ...earnTrainingBernoulliNaiveBayesOpDescSpec.scala | 18 +++++++++ 27 files changed, 236 insertions(+), 5 deletions(-) diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDesc.scala index 3127fa9123..8c5a34c310 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDesc.scala @@ -117,8 +117,12 @@ abstract class SklearnMLOperatorDescriptor[T <: ParamClass] extends PythonOperat | self.dataset = table | | if port == 1 : - | y_train = self.dataset[$groundTruthAttribute] - | X_train = self.dataset[features] + | rows_read = len(self.dataset) + | dataset = self.dataset.dropna(subset=features + [$groundTruthAttribute]) #remove missing values + | if len(dataset) < rows_read: + | print("Skipped", rows_read - len(dataset), "of", rows_read, "rows with missing values") + | y_train = dataset[$groundTruthAttribute] + | X_train = dataset[features] | loop_times = ${getLoopTimes(paraList)} | | for i in range(loop_times): diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnBaggingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnBaggingOpDesc.scala index 59cba35ba6..9581e8f108 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnBaggingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnBaggingOpDesc.scala @@ -22,4 +22,6 @@ package org.apache.texera.amber.operator.sklearn class SklearnBaggingOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.ensemble import BaggingClassifier" override def getUserFriendlyModelName = "Bagging" + // Its default base estimator is a decision tree, which places missing values itself. + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnClassifierOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnClassifierOpDesc.scala index 33c5325b0b..588d717728 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnClassifierOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnClassifierOpDesc.scala @@ -40,8 +40,13 @@ abstract class SklearnClassifierOpDesc extends SklearnModelOpDesc { |class ProcessTableOperator(UDFTableOperator): | @overrides | def process_table(self, table: Table, port: int) -> Iterator[Optional[TableLike]]: + | rows_read = len(table) + | table = $dropMissingRows #remove missing values + | if len(table) < rows_read: + | print("Skipped", rows_read - len(table), "of", rows_read, "rows with missing values") | Y = table[$target] | X = table.drop($target, axis=1) +$reportMissingKept | if port == 0: | self.model = make_pipeline(${vectorizerStage(c => pyb"$c".toString)} ${if ( tfidfTransformer diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDesc.scala index 80827b9c64..fdcd21f8ac 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn class SklearnDecisionTreeOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.tree import DecisionTreeClassifier" override def getUserFriendlyModelName = "Decision Tree" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDummyClassifierOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDummyClassifierOpDesc.scala index 099cf8ce4a..3a7232b405 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDummyClassifierOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnDummyClassifierOpDesc.scala @@ -22,4 +22,7 @@ package org.apache.texera.amber.operator.sklearn class SklearnDummyClassifierOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.dummy import DummyClassifier" override def getUserFriendlyModelName = "Dummy Classifier" + // It predicts from the target's distribution and never reads a feature, so dropping + // a row for a blank feature would change the baseline it is meant to measure. + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreeOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreeOpDesc.scala index c7c5d7d9b8..02787f7030 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreeOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreeOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn class SklearnExtraTreeOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.tree import ExtraTreeClassifier" override def getUserFriendlyModelName = "Extra Tree" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreesOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreesOpDesc.scala index b8bda3b4e7..82af53dd07 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreesOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnExtraTreesOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn class SklearnExtraTreesOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.ensemble import ExtraTreesClassifier" override def getUserFriendlyModelName = "Extra Trees" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDesc.scala index f99da2bff4..e590fa6928 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDesc.scala @@ -53,6 +53,10 @@ class SklearnLinearRegressionOpDesc extends PythonOperatorDescriptor { |class ProcessTableOperator(UDFTableOperator): | @overrides | def process_table(self, table: Table, port: int) -> Iterator[Optional[TableLike]]: + | rows_read = len(table) + | table = table.dropna() #remove missing values + | if len(table) < rows_read: + | print("Skipped", rows_read - len(table), "of", rows_read, "rows with missing values") | Y = table[$target] | X = table.drop($target, axis=1) | if port == 0: diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDesc.scala index fc834659e4..74e474f1ac 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDesc.scala @@ -33,6 +33,7 @@ import com.kjetland.jackson.jsonSchema.annotations.{ } import org.apache.texera.amber.core.tuple.{AttributeType, Schema} import org.apache.texera.amber.pybuilder.PyStringTypes.EncodableString +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.PythonTemplateBuilderStringContext import org.apache.texera.amber.core.workflow.PortIdentity import org.apache.texera.amber.operator.PythonOperatorDescriptor import org.apache.texera.amber.operator.metadata.annotations.{ @@ -127,6 +128,34 @@ abstract class SklearnModelOpDesc extends PythonOperatorDescriptor { @JsonIgnore protected def countVectorizerAlternatives: Option[String] = None + // Tree-based estimators send the missing values of a split down one branch, so a + // blank feature is a signal they can fit on, and the dummy estimator never reads a + // feature at all. Every other estimator here computes over the feature matrix, where + // one NaN spreads through the arithmetic, and scikit-learn refuses the fit rather + // than return a meaningless model. + @JsonIgnore + def handlesMissingValues: Boolean = false + + // A blank target is refused by every estimator, and CountVectorizer calls .lower() + // on each document, so the target and every vectorized column are dropped whatever + // the estimator does. + @JsonIgnore + protected def dropMissingRows: String = + if (countVectorizer) + (text :+ target).map(c => pyb"$c".toString).mkString("table.dropna(subset=[", ", ", "])") + else if (handlesMissingValues) pyb"table.dropna(subset=[$target])".toString + else "table.dropna()" + + // Rows the estimator keeps are still rows the user did not know were incomplete, + // so say how many reached the fit. Empty for the estimators that dropped them all. + @JsonIgnore + protected def reportMissingKept: String = + if (handlesMissingValues && !countVectorizer) + """ | rows_with_gaps = int(X.isna().any(axis=1).sum()) + | | if rows_with_gaps: + | | print("Kept", rows_with_gaps, "rows with missing values, which this model fits without dropping")""".stripMargin + else "" + override def getOutputSchemas( inputSchemas: Map[PortIdentity, Schema] ): Map[PortIdentity, Schema] = { diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDesc.scala index 6e894fccd9..0957339ab4 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDesc.scala @@ -62,9 +62,12 @@ class SklearnPredictionOpDesc extends PythonOperatorDescriptor { | input_features = tuple_ | if $groundTruthAttribute != "": | input_features = input_features.get_partial_tuple([col for col in tuple_.get_field_names() if col != $groundTruthAttribute]) - | tuple_[$resultAttribute] = type(tuple_[$groundTruthAttribute])(self.model.predict(Table.from_tuple_likes([input_features]))[0]) + | if Table.from_tuple_likes([input_features]).isna().any(axis=None): + | tuple_[$resultAttribute] = None #keep the row, leave the result empty | else: - | tuple_[$resultAttribute] = str(self.model.predict(Table.from_tuple_likes([input_features]))[0]) + | prediction = self.model.predict(Table.from_tuple_likes([input_features]))[0] + | #the output schema names this column's type, so reading one off a row could only disagree with it + | tuple_[$resultAttribute] = prediction if $groundTruthAttribute != "" else str(prediction) | yield tuple_""".encode override def operatorInfo: OperatorInfo = diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnRandomForestOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnRandomForestOpDesc.scala index f2f3a51cf8..fd3221a1d2 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnRandomForestOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/SklearnRandomForestOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn class SklearnRandomForestOpDesc extends SklearnClassifierOpDesc { override def getImportStatements = "from sklearn.ensemble import RandomForestClassifier" override def getUserFriendlyModelName = "Random Forest" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDesc.scala index e262bb2953..8e633169ee 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDesc.scala @@ -66,7 +66,12 @@ class SklearnTestingOpDesc extends PythonOperatorDescriptor { | self.data.append(tuple_) | else: | model = tuple_[$model] - | table = Table(self.data) + | #the model arrives already fitted, so this operator cannot ask which + | #estimator it holds and drops on every column to be safe + | rows_read = len(self.data) + | table = Table(self.data).dropna() #remove missing values + | if len(table) < rows_read: + | print("Skipped", rows_read - len(table), "of", rows_read, "rows with missing values") | Y = table[$target] | X = table.drop($target, axis=1) | predictions = model.predict(X.squeeze()) diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBaggingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBaggingOpDesc.scala index 96558a9d97..e82484349c 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBaggingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBaggingOpDesc.scala @@ -22,4 +22,6 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingBaggingOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.ensemble import BaggingClassifier" override def getUserFriendlyModelName = "Training: Bagging" + // Its default base estimator is a decision tree, which places missing values itself. + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDecisionTreeOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDecisionTreeOpDesc.scala index 8fb45fba07..4c3d5409ef 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDecisionTreeOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDecisionTreeOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingDecisionTreeOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.tree import DecisionTreeClassifier" override def getUserFriendlyModelName = "Training: Decision Tree" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDummyClassifierOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDummyClassifierOpDesc.scala index 0423fab054..0f70520a54 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDummyClassifierOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingDummyClassifierOpDesc.scala @@ -22,4 +22,7 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingDummyClassifierOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.dummy import DummyClassifier" override def getUserFriendlyModelName = "Training: Dummy Classifier" + // It predicts from the target's distribution and never reads a feature, so dropping + // a row for a blank feature would change the baseline it is meant to measure. + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreeOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreeOpDesc.scala index 34d362700c..22e8b96013 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreeOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreeOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingExtraTreeOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.tree import ExtraTreeClassifier" override def getUserFriendlyModelName = "Training: Extra Tree" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreesOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreesOpDesc.scala index 155d6facbe..34ced708e4 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreesOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingExtraTreesOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingExtraTreesOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.ensemble import ExtraTreesClassifier" override def getUserFriendlyModelName = "Training: Extra Trees" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingOpDesc.scala index 6d066a5e9b..3de809dd9f 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingOpDesc.scala @@ -40,8 +40,13 @@ class SklearnTrainingOpDesc extends SklearnModelOpDesc { |class ProcessTableOperator(UDFTableOperator): | @overrides | def process_table(self, table: Table, port: int) -> Iterator[Optional[TableLike]]: + | rows_read = len(table) + | table = $dropMissingRows #remove missing values + | if len(table) < rows_read: + | print("Skipped", rows_read - len(table), "of", rows_read, "rows with missing values") | Y = table[$target] | X = table.drop($target, axis=1) +$reportMissingKept | model = make_pipeline(${vectorizerStage(c => pyb"$c".toString)} ${if ( tfidfTransformer ) "TfidfTransformer()," diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingRandomForestOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingRandomForestOpDesc.scala index a9121052d5..ba27315a42 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingRandomForestOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingRandomForestOpDesc.scala @@ -22,4 +22,5 @@ package org.apache.texera.amber.operator.sklearn.training class SklearnTrainingRandomForestOpDesc extends SklearnTrainingOpDesc { override def getImportStatements = "from sklearn.ensemble import RandomForestClassifier" override def getUserFriendlyModelName = "Training: Random Forest" + override def handlesMissingValues = true } diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDescSpec.scala index ba620af298..f151fc3b3f 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/machineLearning/sklearnAdvanced/base/SklearnAdvancedBaseDescSpec.scala @@ -103,6 +103,13 @@ class SklearnAdvancedBaseDescSpec extends AnyFlatSpec with Matchers { code should include("yield df") } + // This family reads a named list of features rather than every column, so the + // drop names those columns: a blank anywhere else must not cost the row. + it should "drop rows missing a selected feature or the ground truth" in { + val d = newOp(List(hyperParam("n_neighbors", "int", fromWorkflow = false, value = "5"))) + d.generatePythonCode() should include("self.dataset.dropna(subset=features + [") + } + it should "loop once when no parameter is sourced from the workflow" in { val d = newOp(List(hyperParam("n_neighbors", "int", fromWorkflow = false, value = "5"))) val code = d.generatePythonCode() diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnBernoulliNaiveBayesOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnBernoulliNaiveBayesOpDescSpec.scala index 0d5fa28fc4..f51e72b26c 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnBernoulliNaiveBayesOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnBernoulliNaiveBayesOpDescSpec.scala @@ -65,6 +65,24 @@ class SklearnBernoulliNaiveBayesOpDescSpec extends AnyFlatSpec with Matchers { code should include("Bernoulli Naive Bayes") } + // The same table statement serves both ports, so training and scoring skip a + // row with a missing value alike. + it should "drop rows with missing values before fitting and before scoring" in { + val d = new SklearnBernoulliNaiveBayesOpDesc + d.target = "y" + d.generatePythonCode() should include("table.dropna()") + } + + // Dropping a row changes the data the model was asked to learn from, so the count + // is reported next to the metrics rather than left for the user to notice. + it should "say how many rows it dropped" in { + val d = new SklearnBernoulliNaiveBayesOpDesc + d.target = "y" + val code = d.generatePythonCode() + code should include("rows_read = len(table)") + code should include("\"Skipped\"") + } + "SklearnBernoulliNaiveBayesOpDesc" should "round-trip its config fields through the polymorphic base" in { val d = new SklearnBernoulliNaiveBayesOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDescSpec.scala index be316adf0f..a2d42832dd 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnDecisionTreeOpDescSpec.scala @@ -47,6 +47,17 @@ class SklearnDecisionTreeOpDescSpec extends AnyFlatSpec with Matchers { d.text shouldBe empty } + // A split compares, and a comparison can ask whether the value is there before it + // asks how large it is, so the row is worth keeping and the count is worth saying. + it should "keep a row with a blank feature and say how many it kept" in { + val d = new SklearnDecisionTreeOpDesc + d.target = "y" + val code = d.generatePythonCode() + code should not include "table.dropna() " + code should include("dropna(subset=[") + code should include("rows_with_gaps") + } + "SklearnDecisionTreeOpDesc.getOutputSchemas" should "emit the model_name/model schema keyed by the declared output port" in { val d = new SklearnDecisionTreeOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDescSpec.scala index 49d4156c05..2c1c02a9eb 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnLinearRegressionOpDescSpec.scala @@ -64,6 +64,16 @@ class SklearnLinearRegressionOpDescSpec extends AnyFlatSpec with Matchers { code should include("class ProcessTableOperator(UDFTableOperator)") } + // This operator builds its own pipeline rather than inheriting the classifier + // base's, so the drop the rest of the family gained has to be stated here too. + it should "drop rows with missing values and say how many" in { + val d = new SklearnLinearRegressionOpDesc + d.target = "y" + val code = d.generatePythonCode() + code should include("table.dropna()") + code should include("\"Skipped\"") + } + "SklearnLinearRegressionOpDesc" should "round-trip its target through the polymorphic base" in { val d = new SklearnLinearRegressionOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDescSpec.scala index 1bf1ca4c4f..127ddabaec 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnModelOpDescSpec.scala @@ -35,6 +35,8 @@ class SklearnModelOpDescSpec extends AnyFlatSpec with Matchers { "from sklearn.linear_model import LogisticRegression" override def getUserFriendlyModelName: String = "Test Model" override def generatePythonCode(): String = "" + // dropMissingRows is protected for the codegen bases; reach it from inside. + def generateDropForTest: String = dropMissingRows override def operatorInfo: OperatorInfo = OperatorInfo( getUserFriendlyModelName, @@ -51,6 +53,43 @@ class SklearnModelOpDescSpec extends AnyFlatSpec with Matchers { d.tfidfTransformer shouldBe false } + // The safe default is the estimator that cannot take a NaN, so an operator only + // keeps incomplete rows when its own class says the estimator places them. + it should "assume an estimator cannot fit a missing value" in { + (new TestSklearnModelOpDesc).handlesMissingValues shouldBe false + } + + "SklearnModelOpDesc.dropMissingRows" should + "drop on every column when the estimator cannot fit a missing value" in { + val d = new TestSklearnModelOpDesc + d.target = "y" + d.generateDropForTest shouldBe "table.dropna()" + } + + // The target is refused by every estimator, so it is dropped even here, but a blank + // feature is left in place for the estimator to make its own use of. + it should "drop on the target alone when the estimator places a missing value" in { + val d = new TestSklearnModelOpDesc { + override def handlesMissingValues = true + } + d.target = "y" + d.generateDropForTest should include("dropna(subset=[") + d.generateDropForTest should not be "table.dropna()" + } + + // CountVectorizer calls .lower() on each document, which a None does not answer, so + // every column it reads goes even for an estimator that would otherwise keep the row. + it should "drop on the text columns too when the vectorizer is on" in { + val d = new TestSklearnModelOpDesc { + override def handlesMissingValues = true + } + d.target = "y" + d.text = List("note", "body") + d.countVectorizer = true + // the two text columns and the target, each named through the decoder + d.generateDropForTest.split("decode_python_template").length - 1 shouldBe 3 + } + "SklearnModelOpDesc.getOutputSchemas" should "key the single output schema by the operator's output port id" in { val d = new TestSklearnModelOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDescSpec.scala index 2b5a76284a..61e6b1750e 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/SklearnPredictionOpDescSpec.scala @@ -92,6 +92,52 @@ class SklearnPredictionOpDescSpec extends AnyFlatSpec with Matchers { code should include("yield tuple_") } + // This operator adds a column to the user's rows, so a row it cannot predict + // on keeps its place with an empty result rather than disappearing. + it should "keep a row with a missing value and leave its result empty" in { + val d = new SklearnPredictionOpDesc + d.model = "model" + d.resultAttribute = "prediction" + val code = d.generatePythonCode() + code should include("isna().any(axis=None)") + code should include("] = None") + } + + // The ignored column is not read by the model, so a blank there must not cost the + // row its prediction: the emptiness test reads the features it actually predicts on. + it should "test the features for emptiness rather than the whole row" in { + val d = new SklearnPredictionOpDesc + d.model = "model" + d.resultAttribute = "prediction" + d.groundTruthAttribute = "y" + val code = d.generatePythonCode() + code should include("Table.from_tuple_likes([input_features]).isna()") + code should not include "Table.from_tuple_likes([tuple_]).isna()" + } + + // The output schema names the result column's type and the framework casts to it, + // so a per-row cast could only disagree with it on the row where the ignored column + // is itself blank and has no type to read off. + it should "not read the result's type off the ignored column" in { + val d = new SklearnPredictionOpDesc + d.model = "model" + d.resultAttribute = "prediction" + d.groundTruthAttribute = "y" + val code = d.generatePythonCode() + code should include("] = prediction if") + code should not include "type(tuple_" + } + + // Without an ignored column the schema declares the result a string, so this is the + // one case where the generated code converts. + it should "write the prediction as text when no ignored column is configured" in { + val d = new SklearnPredictionOpDesc + d.model = "model" + d.resultAttribute = "prediction" + val code = d.generatePythonCode() + code should include("str(prediction)") + } + "SklearnPredictionOpDesc" should "round-trip its config fields through the polymorphic base" in { val d = new SklearnPredictionOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDescSpec.scala index 8c93200500..68ec31a8d7 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/testing/SklearnTestingOpDescSpec.scala @@ -80,6 +80,15 @@ class SklearnTestingOpDescSpec extends AnyFlatSpec with Matchers { code should include(".predict(") } + // The scores are computed over the rows the model can be applied to, the way + // COUNT and MIN are computed over the rows that have a value. + it should "drop rows with missing values before scoring" in { + val d = new SklearnTestingOpDesc + d.model = "model" + d.target = "y" + d.generatePythonCode() should include("Table(self.data).dropna()") + } + "SklearnTestingOpDesc" should "round-trip its config fields through the polymorphic base" in { val d = new SklearnTestingOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBernoulliNaiveBayesOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBernoulliNaiveBayesOpDescSpec.scala index 14da4504db..8bccb3d666 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBernoulliNaiveBayesOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sklearn/training/SklearnTrainingBernoulliNaiveBayesOpDescSpec.scala @@ -64,6 +64,24 @@ class SklearnTrainingBernoulliNaiveBayesOpDescSpec extends AnyFlatSpec with Matc code should include("Training: Bernoulli Naive Bayes") } + // Every column but the target is a feature here, so a row missing any value is + // one the estimator cannot be fitted on. + it should "drop rows with missing values before fitting" in { + val d = new SklearnTrainingBernoulliNaiveBayesOpDesc + d.target = "y" + d.generatePythonCode() should include("table.dropna()") + } + + // With Count Vectorizer on, only the text and target columns are read, so a + // blank in any other column must not cost the row. + it should "drop on the text and target columns only when vectorizing text" in { + val d = new SklearnTrainingBernoulliNaiveBayesOpDesc + d.target = "y" + d.countVectorizer = true + d.text = List("note") + d.generatePythonCode() should include("table.dropna(subset=[") + } + "SklearnTrainingBernoulliNaiveBayesOpDesc" should "round-trip its config fields through the polymorphic base" in { val d = new SklearnTrainingBernoulliNaiveBayesOpDesc d.target = "label"
