diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala index 3853d1cdd67..3fe5e9bec28 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala @@ -21,6 +21,7 @@ package org.apache.texera.amber.operator import org.apache.texera.amber.core.tuple.{AttributeType, Schema} import org.apache.texera.amber.core.workflow.PortIdentity +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral import java.net.URLDecoder import java.nio.charset.StandardCharsets @@ -67,8 +68,20 @@ trait StandaloneCodeGenerator { } /** - * The file's own name, for a script that reads it from its own directory - * rather than through Texera's resolved URI. + * The file this operator reads, as Texera resolved it, or None for one that + * reads no file. + * + * The script cannot open a resolved URI, so the body writes + * [[StandaloneCodeGenerator.SourceFilePlaceholder]] where the file should be + * named and the translator puts a name there. The name has to come from the + * plan: two sources reading different files whose paths end in the same + * segment both asked for `data.csv`, and the script read one of them twice. + */ + def standaloneSourcePath(): Option[String] = None + + /** + * The name to offer for [[standaloneSourcePath]], before the plan has had a + * chance to say whether another source already wants it. * * Taken from the last path segment instead of by parsing the whole string as a * URI: the resolver percent-encodes the file-relative segments but leaves the @@ -76,12 +89,13 @@ trait StandaloneCodeGenerator { * called `v3 - with long text` makes `new URI` throw on the space and no code * is generated at all. */ - protected def sourceBasename(rawPath: String): String = { - val segment = rawPath.split("/").lastOption.getOrElse("") - // Percent-decoding only, matching what `URI.getPath` used to return here: form - // decoding would also turn a literal `+` in a file name into a space. - URLDecoder.decode(segment.replace("+", "%2B"), StandardCharsets.UTF_8) - } + final def standaloneSourceName(): Option[String] = + standaloneSourcePath().map { rawPath => + val segment = rawPath.split("/").lastOption.getOrElse("") + // Percent-decoding only, matching what `URI.getPath` used to return here: form + // decoding would also turn a literal `+` in a file name into a space. + URLDecoder.decode(segment.replace("+", "%2B"), StandardCharsets.UTF_8) + } def producesDataFrame(): Boolean = true @@ -110,3 +124,58 @@ trait StandaloneCodeGenerator { */ def standaloneImports(): Seq[String] = Seq.empty } + +object StandaloneCodeGenerator { + + /** + * What a source writes where the file it reads should be named. + * + * A bare identifier rather than a string literal, because the translator only + * rewrites the code parts of a body and leaves literals and comments alone. + */ + val SourceFilePlaceholder: String = "sourceFile" + + /** Python that gives an empty read the columns and types its schema declares. + * + * With no row there is nothing to infer a type from, so pandas leaves every + * column an object, and a JSONL read of no lines has no columns at all. The + * executor still declares the schema it inferred from the file, so the next + * step looks for those columns in those types. A read with rows is left alone. + */ + def typeAnEmptyRead(frame: String, schema: Schema): String = { + val names = schema.getAttributes.map(a => pyStringLiteral(a.getName)) + val dtypes = schema.getAttributes.flatMap { a => + emptyDtypes.get(a.getType).map(d => s"""${pyStringLiteral(a.getName)}: "$d"""") + } + s"""if $frame.empty: + | $frame = $frame.reindex(columns=[${names.mkString(", ")}]).astype({${dtypes + .mkString(", ")}})""".stripMargin + } + + /** Python that reads each TIMESTAMP column of `schema`, held as text, the way + * the engine reads it, with `_texera_text_to_timestamp` from + * [[StandaloneHelpers.AttributeCasts]], which the caller declares. + * + * pandas parses a date column into nanoseconds, which end in 2262, and left a + * later year as text. The engine hands each field to DateParserUtils, so a + * row states its own format, any year reads, and the reading is cut to the + * millisecond. + */ + def parseTimestamps(frame: String, schema: Schema): String = + schema.getAttributes + .filter(_.getType == AttributeType.TIMESTAMP) + .map { a => + val name = pyStringLiteral(a.getName) + s"$frame[$name] = _texera_text_to_timestamp($frame[$name])" + } + .mkString("\n") + + private val emptyDtypes: Map[AttributeType, String] = Map( + AttributeType.INTEGER -> "Int32", + AttributeType.LONG -> "Int64", + AttributeType.DOUBLE -> "float64", + AttributeType.BOOLEAN -> "boolean", + AttributeType.TIMESTAMP -> "datetime64[us]", + AttributeType.STRING -> "object" + ) +} diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/ScanSourceOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/ScanSourceOpDesc.scala index 20e59f751e7..80458e25e2a 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/ScanSourceOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/ScanSourceOpDesc.scala @@ -82,9 +82,15 @@ abstract class ScanSourceOpDesc extends SourceOperatorDescriptor { /** Rows actually used for type inference: INFER_READ_LIMIT, capped by `windowLimit` * when smaller. + * + * A window of no rows is still a window on this file, and the file's columns do + * not depend on how many of its rows were asked for, so a Limit of 0 samples as + * if no limit were set. Capping by it left nothing to infer from, and the + * operator declared a schema of no columns at all. */ @JsonIgnore - def inferSampleSize: Int = windowLimit.getOrElse(INFER_READ_LIMIT).min(INFER_READ_LIMIT) + def inferSampleSize: Int = + windowLimit.filter(_ > 0).getOrElse(INFER_READ_LIMIT).min(INFER_READ_LIMIT) override def sourceSchema(): Schema = null diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDesc.scala index 1ebbd100402..81c6ec9a277 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDesc.scala @@ -28,15 +28,19 @@ import org.apache.texera.amber.core.tuple.AttributeTypeUtils.inferSchemaFromRows import org.apache.texera.amber.core.tuple.{AttributeType, Schema} import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{PhysicalOp, SchemaPropagationFunc} +import org.apache.texera.amber.operator.{StandaloneCodeGenerator, StandaloneHelpers} +import org.apache.texera.amber.operator.StandaloneCodeGenerator.SourceFilePlaceholder import org.apache.texera.amber.operator.metadata.annotations.UIWidget import org.apache.texera.amber.operator.source.scan.ScanSourceOpDesc import org.apache.texera.amber.operator.source.scan.csv.CSVScanSourceOpExec +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral import org.apache.texera.amber.util.JSONUtils.objectMapper import java.io.{IOException, InputStreamReader} import java.net.URI +import scala.util.Try -class CSVScanSourceOpDesc extends ScanSourceOpDesc { +class CSVScanSourceOpDesc extends ScanSourceOpDesc with StandaloneCodeGenerator { // One character: every reader narrows this with charAt(0), because univocity's // setDelimiter and scala-csv's DefaultCSVFormat both take a Char. @@ -146,4 +150,109 @@ class CSVScanSourceOpDesc extends ScanSourceOpDesc { } + override def standaloneSourcePath(): Option[String] = fileName + + override def standaloneHelpers(): Seq[String] = Seq(StandaloneHelpers.AttributeCasts) + + override def generateStandaloneCode(): String = { + // Resolve the delimiter the same way the parser above does — first character, empty + // means comma — and escape it. Every value the field accepts has to survive this: + // pandas reads a separator longer than one character as a REGULAR EXPRESSION, and a + // backslash spliced raw produced `sep="\"`, which is not valid Python at all. + val sep = customDelimiter.filter(_.nonEmpty).getOrElse(",").charAt(0).toString + // Texera's encoding enum uses values like UTF_8; pandas expects utf-8. + val encoding = fileEncoding.toString.replace("_", "-").toLowerCase + val headerArg = if (hasHeader) "0" else "None" + + val args = scala.collection.mutable.ArrayBuffer[String]() + args += s"filepath_or_buffer=$SourceFilePlaceholder" + args += s"sep=${pyStringLiteral(sep)}" + args += s"""encoding=${pyStringLiteral(encoding)}""" + args += s"header=$headerArg" + + // The parser above sets no null value, so only an empty field is null and every other + // text stands for itself. pandas instead reads a list of words as missing by default, + // "NA" and "null" among them, which turned a column holding the country code NA into + // nulls. Both halves are needed: dropping the default list stops the words, and naming + // the empty string keeps the blank cell null. + args += "keep_default_na=False" + args += """na_values=[""]""" + + // A column holding a null has to be asked for, or pandas widens it to carry + // the hole: an integer through a float, where a LONG past 2^53 comes back + // rounded (9007199254740993 as ...992) and every INTEGER is left a float + // that prints 3 as 3.0, and a boolean to an object column. The nullable + // dtypes carry the hole and keep the type the schema declared. A timestamp + // is read as its text and parsed below, the way the engine parses it. + // + // By position, header or not, because the schema's names are not pandas' + // until the rename below: a blank header is `column-2` here and `Unnamed: 1` + // there, and asking for `column-2` ended the read. pandas takes an integer + // here as a position even where a header spells one. A schema that cannot be + // read (an unresolved file) leaves the argument off rather than failing the + // export. + val nullableDtypes = Map( + AttributeType.INTEGER -> "Int32", + AttributeType.LONG -> "Int64", + AttributeType.BOOLEAN -> "boolean", + AttributeType.TIMESTAMP -> "object" + ) + val typedColumns: Seq[String] = + Try(sourceSchema()).toOption.toSeq.flatMap( + _.getAttributes.zipWithIndex.flatMap { + case (attribute, i) => nullableDtypes.get(attribute.getType).map(d => s"""$i: "$d"""") + } + ) + if (typedColumns.nonEmpty) args += s"dtype={${typedColumns.mkString(", ")}}" + + // Clamped: the property editor refuses a negative, but a plan posted to the API + // can still carry one, and pandas rejects a negative `nrows` outright where the + // executor's `take` simply keeps no rows. + offset.map(_.max(0)).foreach { o => + // With a header, skip offset rows after row 0; without, skip offset rows from the start. + // The end is counted in Long: the largest offset the operator accepts + // overflows an Int on the way past the header, and the range came out + // empty, skipping nothing where the executor's `drop` keeps no rows. Asked + // of each line rather than listed, since pandas makes a set of a list of + // rows to skip, and one of two billion ran the script out of memory. + if (hasHeader) args += s"skiprows=lambda _i: 0 < _i <= ${o.toLong}" + else args += s"skiprows=$o" + } + limit.map(_.max(0)).foreach(l => args += s"nrows=$l") + + // Without a header, an offset at or past the last row left pandas no line to + // count the columns from, and it raised EmptyDataError where the executor + // emits no rows, so the rename and the typing below never ran. Naming the + // columns by position gives it the count and keeps the dtype keys above. + if (!hasHeader) + Try(sourceSchema()).toOption.map(_.getAttributes.size).filter(_ > 0).foreach { n => + args += s"names=[${(0 until n).mkString(", ")}]" + } + + val readCall = s"out1df = pd.read_csv(${args.mkString(", ")})" + + // The schema's own names, which every downstream operator was configured + // against. They differ from pandas' in both directions: a blank header is + // `column-2` here and `Unnamed: 1` there, and a header the user really did + // spell `Unnamed: 1` is kept. Matching the placeholder against the index + // cannot tell those two apart when they coincide; taking the names by + // position can. + val schemaNames: Seq[String] = + Try(sourceSchema()).toOption.toSeq + .flatMap(_.getAttributes.map(a => pyStringLiteral(a.getName))) + + if (schemaNames.nonEmpty) + Seq( + readCall, + s"out1df.columns = [${schemaNames.mkString(", ")}]", + StandaloneCodeGenerator.parseTimestamps("out1df", sourceSchema()), + StandaloneCodeGenerator.typeAnEmptyRead("out1df", sourceSchema()) + ).filter(_.nonEmpty).mkString("\n") + else if (hasHeader) readCall + else { + // Unresolved file: fall back to Texera's headerless naming. + s"""$readCall + |out1df.columns = [f"column-{i + 1}" for i in range(len(out1df.columns))]""".stripMargin + } + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/ParallelCSVScanSourceOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/ParallelCSVScanSourceOpDesc.scala index 4e6da0f5c44..edb667275bf 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/ParallelCSVScanSourceOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csv/ParallelCSVScanSourceOpDesc.scala @@ -29,14 +29,18 @@ import org.apache.texera.amber.core.tuple.AttributeTypeUtils.inferSchemaFromRows import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema} import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{PhysicalOp, SchemaPropagationFunc} +import org.apache.texera.amber.operator.StandaloneCodeGenerator +import org.apache.texera.amber.operator.StandaloneCodeGenerator.SourceFilePlaceholder import org.apache.texera.amber.operator.metadata.annotations.UIWidget import org.apache.texera.amber.operator.source.scan.ScanSourceOpDesc +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral import org.apache.texera.amber.util.JSONUtils.objectMapper import java.io.IOException import java.net.URI +import scala.util.Try -class ParallelCSVScanSourceOpDesc extends ScanSourceOpDesc { +class ParallelCSVScanSourceOpDesc extends ScanSourceOpDesc with StandaloneCodeGenerator { // One character -- see CSVScanSourceOpDesc. @JsonProperty(defaultValue = ",") @@ -81,6 +85,88 @@ class ParallelCSVScanSourceOpDesc extends ScanSourceOpDesc { ) } + override def standaloneSourcePath(): Option[String] = fileName + + override def generateStandaloneCode(): String = { + // First character, empty means comma — the same resolution the reader below does — + // and escaped, so every value the field accepts survives being spliced into Python. + // See CSVScanSourceOpDesc for what handing pandas the raw value did. + val sep = customDelimiter.filter(_.nonEmpty).getOrElse(",").charAt(0).toString + val encoding = fileEncoding.toString.replace("_", "-").toLowerCase + val headerArg = if (hasHeader) "0" else "None" + + val args = scala.collection.mutable.ArrayBuffer[String]() + args += s"filepath_or_buffer=$SourceFilePlaceholder" + args += s"sep=${pyStringLiteral(sep)}" + args += s"""encoding=${pyStringLiteral(encoding)}""" + args += s"header=$headerArg" + + // The block reader nulls an omitted field and leaves every other text alone, + // so only an empty field is missing here. pandas reads a list of words as + // missing by default, "NA" and "null" among them, which turned a column + // holding the country code NA into nulls. See CSVScanSourceOpDesc: both + // halves are needed, one to stop the words and one to keep the blank null. + args += "keep_default_na=False" + args += """na_values=[""]""" + + // Ask for the schema's own type wherever pandas would infer another one. + // The two halves of this operator disagree about a blank: sourceSchema reads + // with scala-csv, where a blank is "" and types its column STRING, while the + // executor nulls it and parses the rest as that STRING. pandas infers a number + // instead, so a column of ids came back as floats, one past 2^53 rounded: + // 9007199254740993 as ...992. A LONG needs the nullable integer for the same + // reason. By position, as in CSVScanSourceOpDesc: under a blank header the + // schema's `column-2` is pandas' `Unnamed: 1`, and pandas passes over a + // name it has no column for, so the type was never applied. + val dtypes: Seq[String] = + Try(sourceSchema()).toOption.toSeq.flatMap( + _.getAttributes.zipWithIndex + .flatMap { + case (a, i) => + val pandasType = a.getType match { + case AttributeType.LONG => Some("Int64") + case AttributeType.STRING => Some("string") + case _ => None + } + pandasType.map(t => s"""$i: "$t"""") + } + ) + if (dtypes.nonEmpty) args += s"dtype={${dtypes.mkString(", ")}}" + + // Limit and offset are inherited fields the parallel reader never reads: + // ParallelCSVScanSourceOpExec.open carves the file into byte ranges and + // leaves both as TODOs. Slicing here gave the export fewer rows than the + // workflow produced, so the window is dropped and said to be dropped. + val ignoredWindow = + if (offset.isEmpty && limit.isEmpty) Seq.empty + else + Seq( + "# NOTE: this operator's limit and offset are ignored, as the parallel CSV reader ignores them." + ) + + val readCall = s"out1df = pd.read_csv(${args.mkString(", ")})" + + // The schema's own names, which every downstream operator was configured + // against: sourceSchema below rewrites a blank header to `column-N`, where + // pandas writes `Unnamed: 1`. Taken by position, so a header the user really + // did spell `Unnamed: 1` survives too. See CSVScanSourceOpDesc. + val schemaNames: Seq[String] = + Try(sourceSchema()).toOption.toSeq + .flatMap(_.getAttributes.map(a => pyStringLiteral(a.getName))) + + val body = + if (schemaNames.nonEmpty) + s"""$readCall + |out1df.columns = [${schemaNames.mkString(", ")}]""".stripMargin + else if (hasHeader) readCall + else + // Unresolved file: fall back to Texera's headerless naming. + s"""$readCall + |out1df.columns = [f"column-{i + 1}" for i in range(len(out1df.columns))]""".stripMargin + + (ignoredWindow :+ body).mkString("\n") + } + override def sourceSchema(): Schema = { val delimiterChar = customDelimiter.filter(_.nonEmpty).getOrElse(",").charAt(0) require( diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDesc.scala index 5eb0b1303cf..b736773ba62 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDesc.scala @@ -28,14 +28,18 @@ import org.apache.texera.amber.core.tuple.AttributeTypeUtils.inferSchemaFromRows import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema} import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{PhysicalOp, SchemaPropagationFunc} +import org.apache.texera.amber.operator.{StandaloneCodeGenerator, StandaloneHelpers} +import org.apache.texera.amber.operator.StandaloneCodeGenerator.SourceFilePlaceholder import org.apache.texera.amber.operator.metadata.annotations.UIWidget import org.apache.texera.amber.operator.source.scan.ScanSourceOpDesc +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral import org.apache.texera.amber.util.JSONUtils.objectMapper import java.io.IOException import java.net.URI +import scala.util.Try -class CSVOldScanSourceOpDesc extends ScanSourceOpDesc { +class CSVOldScanSourceOpDesc extends ScanSourceOpDesc with StandaloneCodeGenerator { // One character -- see CSVScanSourceOpDesc. @JsonProperty(defaultValue = ",") @@ -77,6 +81,89 @@ class CSVOldScanSourceOpDesc extends ScanSourceOpDesc { ) } + override def standaloneSourcePath(): Option[String] = fileName + + override def standaloneHelpers(): Seq[String] = Seq(StandaloneHelpers.AttributeCasts) + + override def generateStandaloneCode(): String = { + // First character, empty means comma — the same resolution the reader below does — + // and escaped, so every value the field accepts survives being spliced into Python. + // See CSVScanSourceOpDesc for what handing pandas the raw value did. + val sep = customDelimiter.filter(_.nonEmpty).getOrElse(",").charAt(0).toString + val encoding = fileEncoding.toString.replace("_", "-").toLowerCase + val headerArg = if (hasHeader) "0" else "None" + + val args = scala.collection.mutable.ArrayBuffer[String]() + args += s"filepath_or_buffer=$SourceFilePlaceholder" + args += s"sep=${pyStringLiteral(sep)}" + args += s"""encoding=${pyStringLiteral(encoding)}""" + args += s"header=$headerArg" + + // This reader has NO missing value at all: scala-csv hands an omitted field + // back as the empty string, and every other text stands for itself, so a + // blank cell is "" and even widens its column to STRING. pandas instead + // reads a list of words as missing by default, "NA" and "null" among them, + // and reads a blank as NaN. Dropping the default list settles both, and no + // na_values is named — where the other CSV readers null a blank, this one + // keeps it. See CSVScanSourceOpDesc. + args += "keep_default_na=False" + + // Read a LONG column as the nullable integer, so a hole does not widen it + // through a float and round the values it carries, and a timestamp as its + // text, parsed below the way the engine parses it. By position, header or + // not, since a blank header is not yet the schema's name. See + // CSVScanSourceOpDesc. + val readAs = Map(AttributeType.LONG -> "Int64", AttributeType.TIMESTAMP -> "object") + val typedColumns: Seq[String] = + Try(sourceSchema()).toOption.toSeq.flatMap( + _.getAttributes.zipWithIndex.flatMap { + case (attribute, i) => readAs.get(attribute.getType).map(d => s"""$i: "$d"""") + } + ) + if (typedColumns.nonEmpty) args += s"dtype={${typedColumns.mkString(", ")}}" + + // Clamped, as in the newer CSV scan: pandas rejects a negative `nrows` where + // the executor's `take` keeps no rows, and only the editor refuses one. + offset.map(_.max(0)).foreach { o => + // Counted in Long past the header and asked of each line, as in the newer + // CSV scan. + if (hasHeader) args += s"skiprows=lambda _i: 0 < _i <= ${o.toLong}" + else args += s"skiprows=$o" + } + limit.map(_.max(0)).foreach(l => args += s"nrows=$l") + + // Without a header, an offset at or past the last row left pandas no line to + // count the columns from, and it raised where the executor emits no rows. + // Named by position, as in the newer CSV scan. + if (!hasHeader) + Try(sourceSchema()).toOption.map(_.getAttributes.size).filter(_ > 0).foreach { n => + args += s"names=[${(0 until n).mkString(", ")}]" + } + + val readCall = s"out1df = pd.read_csv(${args.mkString(", ")})" + + // The schema's own names, which every downstream operator was configured + // against: sourceSchema below rewrites a blank header to `column-N`, where + // pandas writes `Unnamed: 1`. Taken by position, so a header the user really + // did spell `Unnamed: 1` survives too. See CSVScanSourceOpDesc. + val schemaNames: Seq[String] = + Try(sourceSchema()).toOption.toSeq + .flatMap(_.getAttributes.map(a => pyStringLiteral(a.getName))) + + if (schemaNames.nonEmpty) + Seq( + readCall, + s"out1df.columns = [${schemaNames.mkString(", ")}]", + StandaloneCodeGenerator.parseTimestamps("out1df", sourceSchema()), + StandaloneCodeGenerator.typeAnEmptyRead("out1df", sourceSchema()) + ).filter(_.nonEmpty).mkString("\n") + else if (hasHeader) readCall + else + // Unresolved file: fall back to Texera's headerless naming. + s"""$readCall + |out1df.columns = [f"column-{i + 1}" for i in range(len(out1df.columns))]""".stripMargin + } + override def sourceSchema(): Schema = { val delimiterChar = customDelimiter.filter(_.nonEmpty).getOrElse(",").charAt(0) require( @@ -103,16 +190,28 @@ class CSVOldScanSourceOpDesc extends ScanSourceOpDesc { // reopen the file to read from the beginning reader = CSVReader.open(file, fileEncoding.getCharset.name())(CustomFormat) - val startOffset = windowOffset + (if (hasHeader) 1 else 0) - val endOffset = - startOffset + inferSampleSize - val attributeTypeList: Array[AttributeType] = inferSchemaFromRows( - reader.iterator - .slice(startOffset, endOffset) - .map(seq => seq.toArray) - ) - + val header = if (hasHeader) 1 else 0 + // The header and the offset are dropped one after the other, as the executor + // drops them, since their sum is past what an Int holds at the largest offset. + val windowRows = reader.iterator + .drop(header) + .drop(windowOffset) + .take(inferSampleSize) + .map(_.toArray[Any]) + .toSeq reader.close() + // An offset past the last row leaves the window empty, and the types came + // back empty too, so the header below asked a column for a type that was + // never there and the operator threw. The sample is then taken from the + // first row, and a file holding no row at all types every column as text. + val sampleRows = + if (windowRows.nonEmpty) windowRows + else { + val fromStart = CSVReader.open(file, fileEncoding.getCharset.name())(CustomFormat) + try fromStart.iterator.drop(header).take(INFER_READ_LIMIT).map(_.toArray[Any]).toSeq + finally fromStart.close() + } + val attributeTypeList: Array[AttributeType] = inferSchemaFromRows(sampleRows.iterator) // build schema based on inferred AttributeTypes. // Auto-rename blank header positions to `column-N` so empty CSV headers @@ -121,7 +220,7 @@ class CSVOldScanSourceOpDesc extends ScanSourceOpDesc { Schema().add(firstRow.indices.map { i => new Attribute( if (hasHeader && firstRow(i).nonEmpty) firstRow(i) else s"column-${i + 1}", - attributeTypeList(i) + attributeTypeList.lift(i).getOrElse(AttributeType.STRING) ) }) diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpExec.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpExec.scala index 8f6535572ae..f71be829023 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpExec.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpExec.scala @@ -87,8 +87,9 @@ class CSVOldScanSourceOpExec private[csvOld] ( val filePath = DocumentFactory.openReadonlyDocument(new URI(desc.fileName.get)).asFile().toPath reader = CSVReader.open(filePath.toString, desc.fileEncoding.getCharset.name())(CustomFormat) // skip line if this worker reads the start of a file, and the file has a header line - val startOffset = desc.windowOffset + (if (desc.hasHeader) 1 else 0) - rows = reader.iterator.drop(startOffset) + // Dropped one after the other: the largest offset plus the header line is past + // what an Int holds, and the sum wrapped to a negative that skipped no row. + rows = reader.iterator.drop(if (desc.hasHeader) 1 else 0).drop(desc.windowOffset) } override def close(): Unit = { diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDescSpec.scala index f5122381f2a..b5a22dfef70 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csv/CSVScanSourceOpDescSpec.scala @@ -22,6 +22,7 @@ package org.apache.texera.amber.operator.source.scan.csv import com.fasterxml.jackson.databind.JsonNode import com.fasterxml.jackson.databind.node.TextNode import com.github.fge.jsonschema.main.JsonSchemaFactory +import com.typesafe.config.ConfigFactory import org.apache.texera.amber.core.storage.FileResolver import org.apache.texera.amber.core.tuple.{AttributeType, Schema} import org.apache.texera.amber.core.workflow.WorkflowContext.{ @@ -32,14 +33,21 @@ import org.apache.texera.amber.operator.{LogicalOp, TestOperators} import org.apache.texera.amber.operator.metadata.OperatorMetadataGenerator import org.apache.texera.amber.operator.source.scan.ScanSourceOpDesc import org.apache.texera.amber.operator.source.scan.csvOld.CSVOldScanSourceOpDesc -import org.scalatest.BeforeAndAfter +import org.apache.texera.amber.operator.tags.IntegrationTest +import org.apache.texera.amber.util.JSONUtils.objectMapper +import org.scalatest.{BeforeAndAfter, Tag} import org.scalatest.flatspec.AnyFlatSpec import java.nio.charset.StandardCharsets import java.nio.file.Files +import java.util.concurrent.TimeUnit +import scala.io.Source +import scala.util.Try class CSVScanSourceOpDescSpec extends AnyFlatSpec with BeforeAndAfter { + private val NeedsPythonPackages = Tag(classOf[IntegrationTest].getName) + var csvScanSourceOpDesc: CSVScanSourceOpDesc = _ var parallelCsvScanSourceOpDesc: ParallelCSVScanSourceOpDesc = _ before { @@ -97,6 +105,38 @@ class CSVScanSourceOpDescSpec extends AnyFlatSpec with BeforeAndAfter { opDesc.sourceSchema().getAttributes.map(_.getName).toList } + // Writes a CSV whose `big` column holds values past 2^53 and one omitted cell, + // and returns the absolute path. + private def writeNullableLongCsv(): String = { + val tmpFile = Files.createTempFile("nullable-long-", ".csv") + tmpFile.toFile.deleteOnExit() + Files.write( + tmpFile, + "id,big\n1,9007199254740993\n2,\n3,9007199254740995\n".getBytes(StandardCharsets.UTF_8) + ) + tmpFile.toString + } + + // Writes a CSV holding the country code NA and an omitted field, and returns + // the absolute path. + private def writeNaCsv(): String = { + val tmpFile = Files.createTempFile("na-code-", ".csv") + tmpFile.toFile.deleteOnExit() + Files.write(tmpFile, "code,note\nNA,x\n,y\n".getBytes(StandardCharsets.UTF_8)) + tmpFile.toString + } + + // Writes a headered CSV with six data rows and returns the absolute path. + private def writeSixRowCsv(): String = { + val tmpFile = Files.createTempFile("six-row-", ".csv") + tmpFile.toFile.deleteOnExit() + Files.write( + tmpFile, + "id,name\n1,a\n2,b\n3,c\n4,d\n5,e\n6,f\n".getBytes(StandardCharsets.UTF_8) + ) + tmpFile.toString + } + // Writes a numeric column with one blank cell and returns the absolute path. private def writeCsvWithBlankNumericCell(): String = { val tmpFile = Files.createTempFile("blank-cell-", ".csv") @@ -211,6 +251,159 @@ class CSVScanSourceOpDescSpec extends AnyFlatSpec with BeforeAndAfter { ) } + // The name is left to the translator, which is the only thing that can see a + // second source wanting it. What the operator owes is the path and the name it + // would like. + it should "offer the csv basename and read the file by placeholder" in { + csvScanSourceOpDesc.fileName = Some(TestOperators.CountrySalesSmallMultiLineCsvPath) + csvScanSourceOpDesc.customDelimiter = Some(",") + csvScanSourceOpDesc.hasHeader = true + csvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(csvScanSourceOpDesc.fileName.get)) + + val code = csvScanSourceOpDesc.generateStandaloneCode() + + assert(code.contains("filepath_or_buffer=sourceFile")) + assert(csvScanSourceOpDesc.standaloneSourcePath() == csvScanSourceOpDesc.fileName) + assert( + csvScanSourceOpDesc.standaloneSourceName().contains("country_sales_small_multi_line.csv") + ) + assert(!code.contains("base64.b64decode")) + assert(!code.contains("io.BytesIO")) + } + + it should "offer the unresolved csv basename" in { + csvScanSourceOpDesc.fileName = Some(TestOperators.CountrySalesSmallMultiLineCsvPath) + csvScanSourceOpDesc.customDelimiter = Some(",") + csvScanSourceOpDesc.hasHeader = true + + val code = csvScanSourceOpDesc.generateStandaloneCode() + + assert(code.contains("filepath_or_buffer=sourceFile")) + assert( + csvScanSourceOpDesc.standaloneSourceName().contains("country_sales_small_multi_line.csv") + ) + assert(!code.contains("base64.b64decode")) + assert(!code.contains("io.BytesIO")) + } + + // pandas raises nothing for a type asked of a name it has no column for, so + // under a blank header the type was dropped without a word and the column + // was inferred as floats, rounding 9007199254740993 to ...992. + it should "ask pandas for a type under a blank header by its position for parallel CSV" in { + val tmpFile = Files.createTempFile("blank-long-header-", ".csv") + tmpFile.toFile.deleteOnExit() + Files.write( + tmpFile, + "id,\n1,9007199254740993\n2,\n3,9007199254740995\n".getBytes(StandardCharsets.UTF_8) + ) + val path = tmpFile.toString + parallelCsvScanSourceOpDesc.fileName = Some(path) + parallelCsvScanSourceOpDesc.customDelimiter = Some(",") + parallelCsvScanSourceOpDesc.hasHeader = true + parallelCsvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + + assert( + parallelCsvScanSourceOpDesc.sourceSchema().getAttribute("column-2").getType == + AttributeType.STRING + ) + assert(parallelCsvScanSourceOpDesc.generateStandaloneCode().contains("""dtype={1: "string"}""")) + } + + // sourceSchema reads this operator's file with scala-csv, which hands a blank + // back as "", so one blank cell types the whole column STRING — while the + // executor's block reader nulls the blank and parses the rest as that STRING. + // pandas sees the blank as missing and infers a number instead, which read a + // column of ids back as floats and rounded 9007199254740993 to ...992. + it should "read a parallel CSV column the schema typed STRING as text" in { + val path = writeNullableLongCsv() + parallelCsvScanSourceOpDesc.fileName = Some(path) + parallelCsvScanSourceOpDesc.customDelimiter = Some(",") + parallelCsvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + + assert( + parallelCsvScanSourceOpDesc.sourceSchema().getAttribute("big").getType == + AttributeType.STRING + ) + + val exec = new ParallelCSVScanSourceOpExec( + objectMapper.writeValueAsString(parallelCsvScanSourceOpDesc) + ) + exec.open() + val rows = + try exec.produceTuple().map(_.getFields.toList).toList + finally exec.close() + assert(rows.map(_(1)) == List("9007199254740993", null, "9007199254740995")) + + assert( + parallelCsvScanSourceOpDesc + .generateStandaloneCode() + .contains("""dtype={1: "string"}""") + ) + } + + // The block reader nulls an omitted field and leaves every other text alone, so + // "NA" is the country code it says it is. pandas reads it as missing by default, + // so the export read a column of codes as a column of nulls. + it should "read only an empty field as null for parallel CSV, the way its reader does" in { + val path = writeNaCsv() + parallelCsvScanSourceOpDesc.fileName = Some(path) + parallelCsvScanSourceOpDesc.customDelimiter = Some(",") + parallelCsvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + + val exec = new ParallelCSVScanSourceOpExec( + objectMapper.writeValueAsString(parallelCsvScanSourceOpDesc) + ) + exec.open() + val rows = + try exec.produceTuple().map(_.getFields.toList).toList + finally exec.close() + assert(rows == List(List("NA", "x"), List(null, "y"))) + + val code = parallelCsvScanSourceOpDesc.generateStandaloneCode() + assert(code.contains("keep_default_na=False")) + assert(code.contains("""na_values=[""]""")) + } + + it should "give the parallel CSV frame the names the schema gives it" in { + val path = writeCsvWithEmptyHeader() + parallelCsvScanSourceOpDesc.fileName = Some(path) + parallelCsvScanSourceOpDesc.customDelimiter = Some(",") + parallelCsvScanSourceOpDesc.hasHeader = true + parallelCsvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + + assert( + parallelCsvScanSourceOpDesc + .generateStandaloneCode() + .contains("""out1df.columns = ["id", "name", "column-3", "age"]""") + ) + } + + // ParallelCSVScanSourceOpExec.open carves the file into byte ranges and leaves + // limit and offset as TODOs, so the window the panel offers never reaches the + // rows. Slicing in the export handed back fewer rows than the workflow did. + it should "leave limit and offset out of the parallel CSV read, as its reader ignores them" in { + val path = writeSixRowCsv() + parallelCsvScanSourceOpDesc.fileName = Some(path) + parallelCsvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + parallelCsvScanSourceOpDesc.customDelimiter = Some(",") + parallelCsvScanSourceOpDesc.offset = Some(2) + parallelCsvScanSourceOpDesc.limit = Some(2) + + val exec = new ParallelCSVScanSourceOpExec( + objectMapper.writeValueAsString(parallelCsvScanSourceOpDesc) + ) + exec.open() + val rowsRead = + try exec.produceTuple().size + finally exec.close() + assert(rowsRead == 6) + + val code = parallelCsvScanSourceOpDesc.generateStandaloneCode() + assert(!code.contains("skiprows")) + assert(!code.contains("nrows")) + assert(code.startsWith("# NOTE: this operator's limit and offset are ignored")) + } + it should "use comma as the default delimiter when customDelimiter is not set for parallel CSV" in { parallelCsvScanSourceOpDesc.customDelimiter = None @@ -411,4 +604,103 @@ class CSVScanSourceOpDescSpec extends AnyFlatSpec with BeforeAndAfter { assert(columnNames(oldCsv, path) == List("id", "name", "age")) } + // The limit bounded the sample the inference reads as well as the rows the + // operator emits, so a Limit of 0 had nothing to infer from. The three readers + // then failed differently on the same file: these two declared a schema of no + // columns at all, and the old one threw, its header still asking each column + // for a type the empty sample could not give. A file's columns do not depend + // on how many of its rows were asked for. + it should "keep the file's columns when the window asks for no rows" in { + val path = writeSemicolonCsv() + val csv = new CSVScanSourceOpDesc() + csv.customDelimiter = Some(";") + csv.limit = Some(0) + val parallelCsv = new ParallelCSVScanSourceOpDesc() + parallelCsv.customDelimiter = Some(";") + parallelCsv.limit = Some(0) + val oldCsv = new CSVOldScanSourceOpDesc() + oldCsv.customDelimiter = Some(";") + oldCsv.limit = Some(0) + + assert(columnNames(csv, path) == List("id", "name", "age")) + assert(columnNames(parallelCsv, path) == List("id", "name", "age")) + assert(columnNames(oldCsv, path) == List("id", "name", "age")) + } + + // With no header line left after the offset, pandas had nothing to count the + // columns from and the script raised EmptyDataError where the engine emits no + // rows. + it should "read no rows with the schema's columns, as the engine does, when a headerless offset passes the end" taggedAs NeedsPythonPackages in { + val python = resolvePython().getOrElse(cancel("No runnable python executable")) + if (!canImportPandas(python)) cancel(s"'$python' cannot import pandas") + + val tmpFile = Files.createTempFile("headerless-two-row-", ".csv") + tmpFile.toFile.deleteOnExit() + Files.write(tmpFile, "1,a\n2,b\n".getBytes(StandardCharsets.UTF_8)) + val path = tmpFile.toString + csvScanSourceOpDesc.fileName = Some(path) + csvScanSourceOpDesc.customDelimiter = Some(",") + csvScanSourceOpDesc.hasHeader = false + csvScanSourceOpDesc.offset = Some(10) + csvScanSourceOpDesc.setResolvedFileName(FileResolver.resolve(path)) + + val exec = new CSVScanSourceOpExec(objectMapper.writeValueAsString(csvScanSourceOpDesc)) + exec.open() + val engineRows = + try exec.produceTuple().toList + finally exec.close() + assert(engineRows.isEmpty) + + val driver = + s"""import pandas as pd + |${csvScanSourceOpDesc.standaloneImports().mkString("\n")} + |${csvScanSourceOpDesc.standaloneHelpers().mkString("\n")} + |sourceFile = ${objectMapper.writeValueAsString(path)} + |${csvScanSourceOpDesc.generateStandaloneCode()} + |print(len(out1df), list(out1df.columns)) + |""".stripMargin + val script = Files.createTempFile("csv-standalone-", ".py") + script.toFile.deleteOnExit() + Files.write(script, driver.getBytes(StandardCharsets.UTF_8)) + val process = new ProcessBuilder(python, script.toString).redirectErrorStream(true).start() + val out = Source.fromInputStream(process.getInputStream).mkString + process.waitFor(120, TimeUnit.SECONDS) + + withClue(s"python said:\n$out\nscript:\n$driver") { + assert(process.exitValue() == 0) + assert(out.trim == "0 ['column-1', 'column-2']") + assert( + csvScanSourceOpDesc.sourceSchema().getAttributes.map(_.getName).toList == + List("column-1", "column-2") + ) + } + } + + private def resolvePython(): Option[String] = { + def fromConfig: Option[String] = + Try(ConfigFactory.parseResources("udf.conf").resolve()).toOption + .orElse(Try(ConfigFactory.load()).toOption) + .flatMap(c => Try(c.getConfig("python").getString("path")).toOption) + .map(_.trim) + .filter(_.nonEmpty) + + def runnable(exe: String): Boolean = + Try(new ProcessBuilder(exe, "--version").redirectErrorStream(true).start()).toOption + .exists { p => + if (!p.waitFor(5, TimeUnit.SECONDS)) { p.destroyForcibly(); false } + else p.exitValue() == 0 + } + + (fromConfig.toList ++ List("python3", "python", "py")).distinct.find(runnable) + } + + private def canImportPandas(python: String): Boolean = + Try( + new ProcessBuilder(python, "-c", "import pandas").redirectErrorStream(true).start() + ).toOption + .exists { p => + if (!p.waitFor(60, TimeUnit.SECONDS)) { p.destroyForcibly(); false } + else p.exitValue() == 0 + } + } diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDescSpec.scala index b72476641b0..2dc214aaa04 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/scan/csvOld/CSVOldScanSourceOpDescSpec.scala @@ -19,17 +19,30 @@ package org.apache.texera.amber.operator.source.scan.csvOld +import com.typesafe.config.ConfigFactory import org.apache.texera.amber.core.executor.OpExecWithClassName +import org.apache.texera.amber.core.storage.FileResolver +import org.apache.texera.amber.core.tuple.AttributeType import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.operator.LogicalOp import org.apache.texera.amber.operator.metadata.OperatorGroupConstants import org.apache.texera.amber.operator.source.scan.FileDecodingMethod +import org.apache.texera.amber.operator.tags.IntegrationTest import org.apache.texera.amber.util.JSONUtils.objectMapper +import org.scalatest.Tag import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers +import java.nio.charset.StandardCharsets +import java.nio.file.Files +import java.util.concurrent.TimeUnit +import scala.io.Source +import scala.util.Try + class CSVOldScanSourceOpDescSpec extends AnyFlatSpec with Matchers { + private val NeedsPythonPackages = Tag(classOf[IntegrationTest].getName) + private val workflowId = WorkflowIdentity(1L) private val executionId = ExecutionIdentity(1L) @@ -95,4 +108,103 @@ class CSVOldScanSourceOpDescSpec extends AnyFlatSpec with Matchers { r.limit shouldBe Some(10) r.offset shouldBe Some(5) } + + // The limit bounded the sample the inference reads as well as the rows the + // operator emits, so a Limit of 0 had nothing to infer from. The types came + // back empty while the header still asked each column for one, and the + // operator threw before a row was read. A file's columns do not depend on how + // many of its rows were asked for. + "CSVOldScanSourceOpDesc.sourceSchema" should "keep the file's columns when the window asks for no rows" in { + val d = describing(writeCsv("id,name\n1,alice\n2,bob\n")) + d.limit = Some(0) + + val schema = d.sourceSchema() + schema.getAttributeNames shouldBe List("id", "name") + schema.getAttribute("id").getType shouldBe AttributeType.INTEGER + schema.getAttribute("name").getType shouldBe AttributeType.STRING + } + + // With no header line left after the offset, pandas had nothing to count the + // columns from and the script raised where the engine emits no rows. + "CSVOldScanSourceOpDesc.generateStandaloneCode" should + "read no rows with the schema's columns, as the engine does, when a headerless offset passes the end" taggedAs NeedsPythonPackages in { + val python = resolvePython().getOrElse(cancel("No runnable python executable")) + if (!canImportPandas(python)) cancel(s"'$python' cannot import pandas") + + val path = writeCsv("1,a\n2,b\n") + val d = describing(path) + d.hasHeader = false + d.offset = Some(10) + + val exec = new CSVOldScanSourceOpExec(objectMapper.writeValueAsString(d)) + exec.open() + val engineRows = + try exec.produceTuple().toList + finally exec.close() + engineRows shouldBe empty + + val driver = + s"""import pandas as pd + |${d.standaloneImports().mkString("\n")} + |${d.standaloneHelpers().mkString("\n")} + |sourceFile = ${objectMapper.writeValueAsString(path)} + |${d.generateStandaloneCode()} + |print(len(out1df), list(out1df.columns)) + |""".stripMargin + val script = Files.createTempFile("csv-old-standalone-", ".py") + script.toFile.deleteOnExit() + Files.write(script, driver.getBytes(StandardCharsets.UTF_8)) + val process = new ProcessBuilder(python, script.toString).redirectErrorStream(true).start() + val out = Source.fromInputStream(process.getInputStream).mkString + process.waitFor(120, TimeUnit.SECONDS) + + withClue(s"python said:\n$out\nscript:\n$driver") { + process.exitValue() shouldBe 0 + out.trim shouldBe "0 ['column-1', 'column-2']" + d.sourceSchema().getAttributeNames shouldBe List("column-1", "column-2") + } + } + + private def resolvePython(): Option[String] = { + def fromConfig: Option[String] = + Try(ConfigFactory.parseResources("udf.conf").resolve()).toOption + .orElse(Try(ConfigFactory.load()).toOption) + .flatMap(c => Try(c.getConfig("python").getString("path")).toOption) + .map(_.trim) + .filter(_.nonEmpty) + + def runnable(exe: String): Boolean = + Try(new ProcessBuilder(exe, "--version").redirectErrorStream(true).start()).toOption + .exists { p => + if (!p.waitFor(5, TimeUnit.SECONDS)) { p.destroyForcibly(); false } + else p.exitValue() == 0 + } + + (fromConfig.toList ++ List("python3", "python", "py")).distinct.find(runnable) + } + + private def canImportPandas(python: String): Boolean = + Try( + new ProcessBuilder(python, "-c", "import pandas").redirectErrorStream(true).start() + ).toOption + .exists { p => + if (!p.waitFor(60, TimeUnit.SECONDS)) { p.destroyForcibly(); false } + else p.exitValue() == 0 + } + + private def writeCsv(content: String): String = { + val file = Files.createTempFile("csv-old-", ".csv") + file.toFile.deleteOnExit() + Files.write(file, content.getBytes(StandardCharsets.UTF_8)) + file.toString + } + + private def describing(path: String): CSVOldScanSourceOpDesc = { + val d = new CSVOldScanSourceOpDesc + d.fileName = Some(path) + d.customDelimiter = Some(",") + d.hasHeader = true + d.setResolvedFileName(FileResolver.resolve(path)) + d + } } diff --git a/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala b/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala index 8e873403787..7fa1be9c752 100644 --- a/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala +++ b/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala @@ -25,7 +25,9 @@ import org.apache.texera.amber.core.virtualidentity.OperatorIdentity import org.apache.texera.amber.core.workflow.PortIdentity import org.apache.texera.common.compiler.model.LogicalPlan import org.apache.texera.amber.operator.StandaloneCodeGenerator +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral +import java.util.regex.Matcher import scala.collection.mutable import scala.collection.mutable.ArrayBuffer import scala.jdk.CollectionConverters._ @@ -67,6 +69,25 @@ class WorkflowToPythonTranslator extends LazyLogging { // getTopologicalOpIds() uses jgrapht internally — no need for a custom topo sort val topoOrder = logicalPlan.getTopologicalOpIds.asScala.toList + // What each source's file is called in the directory the script runs from. + // A source offers the last segment of its resolved path, which is the name a + // person would give the file, but two sources reading different files can + // offer the same one — and did, leaving the script to read one of them twice + // and say nothing. The second gets a name of its own, for the same reason a + // chart's output file is numbered. Keyed by the resolved path, so one file + // read by two operators keeps one name. + val sourceFileNames = mutable.Map[String, String]() + topoOrder.map(logicalPlan.getOperator).foreach { + case gen: StandaloneCodeGenerator => + gen.standaloneSourcePath().foreach { path => + sourceFileNames.getOrElseUpdate( + path, + distinctName(gen.standaloneSourceName().getOrElse(""), sourceFileNames.values.toSet) + ) + } + case _ => () + } + // pandas is the one module every generator uses: an operator body reads and // writes frames whatever else it does. Everything beyond that is asked of the // operators in the plan, so a script that draws nothing does not require a @@ -147,6 +168,7 @@ class WorkflowToPythonTranslator extends LazyLogging { inVars, outVars, fileBase(displayName, fileBaseCounts), + gen.standaloneSourcePath().flatMap(sourceFileNames.get).getOrElse(""), displayName ) @@ -214,6 +236,18 @@ class WorkflowToPythonTranslator extends LazyLogging { s"${stem}_$n" } + // The offered name if no other source has taken it, otherwise the same name + // numbered before its extension: data.csv, then data-2.csv. Numbering the stem + // rather than appending keeps the suffix, which is what a reader opens the + // file by. + private def distinctName(offered: String, taken: Set[String]): String = { + if (!taken.contains(offered)) return offered + val dot = offered.lastIndexOf('.') + val (stem, ext) = + if (dot <= 0) (offered, "") else (offered.substring(0, dot), offered.substring(dot)) + Iterator.from(2).map(n => s"$stem-$n$ext").find(!taken.contains(_)).get + } + // Replaces in{N}df / out{N}df placeholders with concrete variable names. // Substitutes in reverse index order to prevent partial matches (e.g. in1df // inside in10df). Only the code parts are rewritten: a generator writes a @@ -226,6 +260,7 @@ class WorkflowToPythonTranslator extends LazyLogging { inVars: List[String], outVars: List[String], fileBase: String, + sourceFile: String, displayName: String ): String = { def substitute(fragment: String): String = { @@ -238,6 +273,15 @@ class WorkflowToPythonTranslator extends LazyLogging { result = result.replaceAll("""\boutputHtml\b""", "\"" + fileBase + ".html\"") result = result.replaceAll("""\boutputJson\b""", "\"" + fileBase + ".json\"") + // A source names the file it reads sourceFile and gets back the name + // assigned where sourceFileNames is built. Quoted through the escaper the + // operators use, and then quoted again for the replacement: a file name is + // the user's text, so it can hold both a backslash and a `$`. + result = result.replaceAll( + s"""\\b${StandaloneCodeGenerator.SourceFilePlaceholder}\\b""", + Matcher.quoteReplacement(pyStringLiteral(sourceFile)) + ) + // A variadic port takes as many upstream links as the user draws, and an // operator reading one cannot name them: `in1df`/`in2df` state a count, and // whichever count it states is wrong for every other workflow. This one diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala index ba0c3f3ae70..b1321a3bc68 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala @@ -163,6 +163,47 @@ class WorkflowToPythonTranslatorSpec extends AnyFlatSpec with Matchers { script should include("""fig.write_html("stub_2.html")""") } + /** A source offers the last segment of its resolved path, which two sources + * reading different files can spell the same. Handing both that one name left + * the script reading one file twice and saying nothing about it. + */ + it should "give each source reading a different file a name of its own" in { + val script = translateSources("/alice/sales/v1/data.csv", "/bob/ops/v3/data.csv") + script should include("""pd.read_csv("data.csv")""") + script should include("""pd.read_csv("data-2.csv")""") + } + + /** The other half: one file read twice is still one file, so numbering it would + * send the reader looking for a second copy that was never there. + */ + it should "give two sources reading one file the same name" in { + val script = translateSources("/alice/sales/v1/data.csv", "/alice/sales/v1/data.csv") + script.linesIterator.count(_.contains("""pd.read_csv("data.csv")""")) shouldBe 2 + script should not include "data-2.csv" + } + + /** A file name is the user's text, so it reaches the script through the escaper + * the operators use and through the replacement quoting on top of that: a bare + * `$` in a name is a group reference to `replaceAll`. + */ + it should "escape a source file name that Python or the replacement would read" in { + val script = translateSources("""/alice/v1/we"ird$1\x.csv""") + script should include("""pd.read_csv("we\"ird$1\\x.csv")""") + } + + private def translateSources(paths: String*): String = { + val ops = paths.zipWithIndex.map { + case (path, i) => + val op = + new StubOp(s"out1df = pd.read_csv(${StandaloneCodeGenerator.SourceFilePlaceholder})") { + override def standaloneSourcePath(): Option[String] = Some(path) + } + op.setOperatorId(s"source$i") + op + } + new WorkflowToPythonTranslator().translate(LogicalPlan(ops.toList, List.empty)) + } + /** The translator's own contract when it meets an operator it cannot render: * a comment rather than a silently wrong line. */