Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
46 commits
Select commit Hold shift + click to select a range
5b8dd57
feat(workflow-compiling-service): export a workflow as a standalone P…
kz930 Sep 1, 2026
a7558eb
test(workflow-compiling-service): run an operator both ways and compa…
kz930 Sep 1, 2026
cba725f
test(workflow-compiling-service): verify a generated script against t…
kz930 Sep 1, 2026
97e05ec
Merge remote-tracking branch 'upstream/main' into feat/standalone-sou…
kz930 Sep 2, 2026
41724ff
feat(workflow-operator): export the source operators as Python
kz930 Sep 2, 2026
bc8fc17
ci: give the verify spec a job with an interpreter, and keep it out o…
kz930 Sep 2, 2026
9be0dbd
Merge remote-tracking branch 'myfork/feat/standalone-verify-harness' …
kz930 Sep 2, 2026
883ea72
test(verify): pin the no-generator case to an operator that can never…
kz930 Sep 2, 2026
5a44c3c
chore: leave the harness to the change that introduces it
kz930 Sep 2, 2026
3f07020
Merge upstream/main
kz930 Sep 2, 2026
df10bd8
chore: leave the two verify files to #8327 as well
kz930 Sep 2, 2026
73375b9
fix(operator): slice the raw lines before converting them
kz930 Sep 9, 2026
cb23c1e
fix(operator): read a CSV the way the parser does, not the way pandas…
kz930 Sep 10, 2026
c2e2eb3
fix(operator): take the row's first string field as the file to scan
kz930 Sep 10, 2026
1c0405f
fix(operator): decode a fetched body by replacement, as the executor …
kz930 Sep 10, 2026
208816d
fix(operator): slice a JSONL file's lines before parsing them
kz930 Sep 10, 2026
566aa15
fix(operator): take a CSV's column names from the schema
kz930 Sep 10, 2026
3d12fae
test(operator): assert the column names the CSV reader now writes
kz930 Sep 10, 2026
14b5571
Merge remote-tracking branch 'upstream/main' into HEAD
kz930 Sep 11, 2026
832b914
Merge branch 'main' into feat/standalone-sources
kz930 Sep 15, 2026
edf46a4
fix(operator): read an Arrow file into the dtypes that keep its nulls
kz930 Sep 18, 2026
fe5e8c0
test(operator): pin that an Arrow read keeps a missing value apart fr…
kz930 Sep 18, 2026
32173ef
fix(workflow-operator): follow the native readers in the exported sou…
kz930 Sep 18, 2026
292ddbe
Merge remote-tracking branch 'myfork/feat/standalone-sources' into fe…
kz930 Sep 18, 2026
d938d35
feat(workflow-operator): clamp the exported scan window at zero
kz930 Sep 18, 2026
a2290cd
fix(workflow-operator): give a flattened array the executor's column …
kz930 Sep 18, 2026
2ef4c15
fix(operator): order the exported JSONL read the way the schema does
kz930 Sep 19, 2026
46e62a8
fix(operator): read an Arrow file's narrow numeric widths on both paths
kz930 Sep 19, 2026
b310855
fix(operator): read an Arrow file's unsigned integers as the numbers …
kz930 Sep 19, 2026
3700e19
fix(workflow-operator): count the end of an exported scan window in Long
kz930 Sep 19, 2026
d8a9d13
fix(operator): decode a file scan with the charset its Encoding field…
kz930 Sep 19, 2026
95ab40c
fix(operator): name the file each line came from when a file scan was…
kz930 Sep 19, 2026
f3d0db2
fix(operator): type a flattened JSONL timestamp as the schema declare…
kz930 Sep 19, 2026
3d1a4be
fix(operator): keep a file's columns when its scan window asks for no…
kz930 Sep 19, 2026
ce52ee2
fix(operator): read every column an Arrow file states, index or not
kz930 Sep 21, 2026
ce2e484
fix(operator): read an Arrow file's columns under the names it states
kz930 Sep 21, 2026
719cb93
fix(operator): read a zoned Arrow timestamp as the wall clock it holds
kz930 Sep 23, 2026
5bc1a0d
fix(operator): ask pandas for a CSV column's type by its position
kz930 Sep 23, 2026
053a6fd
fix(operator): keep a JSONL boolean with a missing key boolean
kz930 Sep 23, 2026
6d36af8
Merge upstream/main into feat/standalone-sources
kz930 Sep 24, 2026
07798b7
fix(operator): read what the executor reads in the exported sources
kz930 Sep 25, 2026
f874af5
fix(operator): read timestamps as the engine does, and hand Python a …
kz930 Sep 27, 2026
7ed4f40
refactor(operator): leave the JSONL, Arrow and line sources to the tw…
kz930 Sep 27, 2026
6b0a3b5
Merge upstream/main into feat/standalone-sources
kz930 Oct 1, 2026
37f8cc2
Merge upstream/main into feat/standalone-sources
kz930 Oct 5, 2026
280aaeb
fix(workflow-operator): read no rows when a headerless CSV offset pas…
kz930 Oct 8, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -67,21 +68,34 @@ 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
* repository and version names as the user typed them, so a dataset version
* 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

Expand Down Expand Up @@ -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"
)
}
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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(", ")})"
Comment thread
kz930 marked this conversation as resolved.

// 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
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 = ",")
Expand Down Expand Up @@ -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(
Expand Down
Loading
Loading