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-8356-5042d96ec85d98ed18bde841b2e74922c43f5c3b in repository https://gitbox.apache.org/repos/asf/texera.git
commit 0449b467b0fbf0314dd307dac1a5b4f61d5b6441 Author: Kary Zheng <[email protected]> AuthorDate: Mon Sep 14 19:26:45 2026 +0000 test(verify): run an operator through the engine and keep what it wrote (#8356) ### What changes were proposed in this PR? An operator's answer cannot be compared against anything until there is a way to get one. `OpExecHarness` builds the physical operator a descriptor describes, feeds it the rows of a JSONL file per input port, and writes what each output port produced back out, schema in a sidecar because JSONL carries values alone and cannot say a column is an integer rather than a number. The tests that need a Python interpreter are tagged and split into a job that provisions one, so the job that does not stays as fast as it was. ### Any related issues, documentation, discussions? Part of #8325, 3 of 27; that issue lists the set in order. Closes #8408, the task this change is the whole of. ### How was this PR tested? The tests in this change cover it. The whole set is exercised together once the last piece lands: every operator run through the engine and through its generated script, and the two answers compared. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Opus 5) --------- Co-authored-by: Claude Opus 5 (1M context) <[email protected]> --- .github/workflows/build.yml | 65 +++ workflow-compiling-service/build.sbt | 9 + .../translator/verify/tags/IntegrationTest.java | 56 +++ .../amber/translator/verify/OpExecHarness.scala | 441 +++++++++++++++++++++ 4 files changed, 571 insertions(+) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index c0107997b3..0614d54ae0 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -783,6 +783,13 @@ jobs: env: JAVA_OPTS: -Xms2048M -Xmx2048M -Xss6M -XX:ReservedCodeCacheSize=256M -Dfile.encoding=UTF-8 JVM_OPTS: -Xms2048M -Xmx2048M -Xss6M -XX:ReservedCodeCacheSize=256M -Dfile.encoding=UTF-8 + # Exclude @IntegrationTest-tagged specs (workflow-compiling-service's + # OperatorBehaviorSpec forks Python, which this job does not provision). + # Those run in the platform-integration job's workflow-compiling-service + # matrix entry, which provisions Python. Read only by + # workflow-compiling-service/build.sbt; a no-op for the other services + # in this matrix. + WCS_TEST_FILTER: skip-integration services: # Each platform service transitively depends on DAO, which runs JOOQ # code generation at compile time and needs the live texera schema. @@ -972,6 +979,45 @@ jobs: - uses: coursier/cache-action@95e5b1029b6b86e7bac033ee44a0697d8a527d2d # v8.1.1 with: extraSbtFiles: '["*.sbt", "project/**.{scala,sbt}", "project/build.properties" ]' + # --- workflow-compiling-service only: provision Python so the verify + # spec (OperatorBehaviorSpec) can fork real interpreters. These four + # steps mirror the retired standalone workflow-compiling-service-integration + # job; they are no-ops for every other service in the matrix. + - name: Setup Python for Scala-Python verification tests + if: ${{ matrix.service == 'workflow-compiling-service' }} + uses: actions/setup-python@v7 + with: + python-version: "3.12" + - name: Install Python dependencies + # OperatorBehaviorSpec forks Python subprocesses that import pandas / + # pyarrow / plotly and pytexera (the harness adds amber/src/main/python + # to PYTHONPATH). Install amber's runtime deps, same as amber-integration. + # dev-requirements.txt provides the betterproto plugin used by + # bin/python-proto-gen.sh in the proto-generation step below. + if: ${{ matrix.service == 'workflow-compiling-service' }} + run: | + python -m pip install uv + if [ -f amber/requirements.txt ]; then uv pip install --system --index-strategy unsafe-best-match -r amber/requirements.txt; fi + if [ -f amber/operator-requirements.txt ]; then uv pip install --system --index-strategy unsafe-best-match -r amber/operator-requirements.txt; fi + if [ -f amber/dev-requirements.txt ]; then uv pip install --system --index-strategy unsafe-best-match -r amber/dev-requirements.txt; fi + - name: Install protoc + # Path A forks py_op_driver, which imports pyamber and hence the + # generated betterproto bindings in amber/src/main/python/proto (gitignored, + # not checked in). Pin protoc to bin/protoc-version.txt via the upstream + # release zip. Linux-only: this job runs on ubuntu-latest. + if: ${{ matrix.service == 'workflow-compiling-service' }} + run: | + PROTOC_VERSION=$(cat bin/protoc-version.txt) + curl -fsSL -o /tmp/protoc.zip "https://github.com/protocolbuffers/protobuf/releases/download/v${PROTOC_VERSION}/protoc-${PROTOC_VERSION}-linux-x86_64.zip" + sudo unzip -o /tmp/protoc.zip -d /usr/local + sudo chmod +x /usr/local/bin/protoc + sudo chmod -R a+rX /usr/local/include/google + - name: Generate Python proto bindings + # Regenerate amber/src/main/python/proto so the forked py_op_driver can + # import pyamber; without this Path A fails with ImportError on proto + # symbols (e.g. ChannelIdentity). + if: ${{ matrix.service == 'workflow-compiling-service' }} + run: bash bin/python-proto-gen.sh - name: Create Databases run: | psql -h localhost -U postgres -f sql/texera_ddl.sql @@ -1053,6 +1099,25 @@ jobs: # smoke-boot's verdict is LISTEN-based, never log-scraping (#6332). TEXERA_SERVICE_LOG_LEVEL: ${{ runner.debug == '1' && 'DEBUG' || 'WARN' }} run: .github/scripts/smoke-boot.sh "/tmp/dists/${{ matrix.service }}-*/bin/${{ matrix.service }}" "${{ matrix.port }}" + - name: Run workflow-compiling-service Python e2e (integration) tests + # workflow-compiling-service only: the @IntegrationTest-tagged verify + # spec (OperatorBehaviorSpec) forks real Python processes and compares + # translator-generated standalone code against platform output, operator + # by operator. Reuses this job's compiled WCS (the dist step above) and + # postgres instead of a separate integration job. + # WCS_TEST_FILTER=integration-only keeps only @IntegrationTest specs and + # bounds ScalaTest's parallel pool to 4 threads (build.sbt). Every operator + # carrying standalone code runs: the spec's narrowing knobs + # (VERIFY_ONLY / VERIFY_SKIP) are unset here, as they are in any run that + # does not ask for less, and what stays withheld is a single variant or an + # operator that cannot run, recorded against an issue in the runner itself. + # UDF_PYTHON_PATH points the harness at the provisioned interpreter (bare + # name resolves via PATH); it auto-locates amber/src/main/python. + if: ${{ matrix.service == 'workflow-compiling-service' }} + env: + WCS_TEST_FILTER: integration-only + UDF_PYTHON_PATH: python + run: sbt "WorkflowCompilingService/test" pyamber: if: ${{ inputs.run_pyamber }} diff --git a/workflow-compiling-service/build.sbt b/workflow-compiling-service/build.sbt index 1c440d14b0..2427f479a1 100644 --- a/workflow-compiling-service/build.sbt +++ b/workflow-compiling-service/build.sbt @@ -46,6 +46,15 @@ ThisBuild / conflictManager := ConflictManager.latestRevision // tests *within* a suite (e.g. OperatorBehaviorSpec) via ScalaTest's own pool. Global / concurrentRestrictions += Tags.limit(Tags.Test, 1) +// The fast-unit / integration test split; the selection logic itself is shared +// in project/TestFilters.scala. The tag it names is the one this change adds, +// so the two arrive together and the filter never selects on a tag nothing +// carries. +Test / testOptions ++= TestFilters.integrationSplit( + envVar = "WCS_TEST_FILTER", + tag = "org.apache.texera.amber.translator.verify.tags.IntegrationTest" +) + // -P4 bounds ScalaTest's ParallelTestExecution pool, and only this module wants // it: OperatorBehaviorSpec forks a Python subprocess per operator, and at // core-count concurrency (e.g. 12) resource contention caused rare flakes. A diff --git a/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java b/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java new file mode 100644 index 0000000000..4da3aa9cd2 --- /dev/null +++ b/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.texera.amber.translator.verify.tags; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import org.scalatest.TagAnnotation; + +/** + * Class-level marker tag for workflow-compiling-service ScalaTest specs that + * exercise both Scala and Python end-to-end (they fork a real Python process + * to run and compare the translator-generated code). Routing to the + * {@code workflow-compiling-service-integration} CI job is by ScalaTest tag + * filtering, controlled by the {@code WCS_TEST_FILTER} env var in + * {@code workflow-compiling-service/build.sbt}: the lighter + * {@code workflow-compiling-service} test run uses {@code skip-integration} + * (which passes {@code -l org.apache.texera.amber.translator.verify.tags.IntegrationTest} + * to ScalaTest), and the integration job uses {@code integration-only} (which + * passes {@code -n} for the same tag). + * + * <p>Mirrors amber's {@code org.apache.texera.amber.tags.IntegrationTest}. Only + * {@code OperatorBehaviorSpec} carries this tag today — it is the sole spec that + * spawns Python; the other verify specs only exercise pure-JVM classification + * and comparison logic. + * + * <p>Written in Java rather than Scala because ScalaTest detects tag + * annotations via {@code java.lang.annotation} reflection. A Scala + * {@code class extends StaticAnnotation} does not produce a JVM annotation + * interface that {@code @TagAnnotation} can attach to, so the tag would be + * invisible to ScalaTest at runtime. + */ +@TagAnnotation +@Retention(RetentionPolicy.RUNTIME) +@Target({ElementType.METHOD, ElementType.TYPE}) +public @interface IntegrationTest { +} diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala new file mode 100644 index 0000000000..01e48e40da --- /dev/null +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala @@ -0,0 +1,441 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.texera.amber.translator.verify + +import com.fasterxml.jackson.databind.node.ObjectNode +import com.typesafe.scalalogging.LazyLogging +import org.apache.texera.amber.core.executor.{ExecFactory, OpExecWithClassName, OperatorExecutor} +import org.apache.texera.amber.core.tuple.{AttributeType, Schema, Tuple, TupleLike} +import org.apache.texera.amber.core.virtualidentity.{ + ExecutionIdentity, + PhysicalOpIdentity, + WorkflowIdentity +} +import org.apache.texera.amber.core.workflow.{PhysicalOp, PhysicalPlan, PortIdentity} +import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.util.JSONUtils.objectMapper + +import java.nio.file.{Files, Path} +import java.sql.Timestamp +import java.util.Base64 +import scala.collection.mutable +import scala.jdk.CollectionConverters._ + +/** + * Generic harness that drives an OpDesc's OpExec(s) directly, bypassing the + * Pekko/actor runtime. + * + * The wiring, meaning how many OpExecs there are, the internal links and the + * input-port dependency order, comes from `opDesc.getPhysicalPlan(...)`, so a + * join's build and probe, a split's two outputs and a source's empty input map + * all work without a harness change. I/O is JSON Lines with a + * `*.jsonl.schema.json` sidecar per file. + * + * It runs one worker, idx=0 of 1, so nothing here coordinates partitioners + * across executors. Python UDFs go to [[PyOpExecHarness]], which has a real + * worker to give them. + */ +object OpExecHarness extends LazyLogging { + + // Test-only workflow / execution IDs. The values don't matter — the harness + // never persists state under them — but the PhysicalOp factory needs *some* + // IDs to embed in PhysicalOpIdentity. + private val TestWorkflowId = WorkflowIdentity(0L) + private val TestExecutionId = ExecutionIdentity(0L) + + /** + * @param outputs external output port → JSONL file path + * @param outputSchemas same keys as outputs, gives each port's [[Schema]] + */ + final case class Result( + outputs: Map[PortIdentity, Path], + outputSchemas: Map[PortIdentity, Schema] + ) + + /** + * Run `opDesc` against the given inputs and write the outputs to `outputDir`. + * + * @param inputs map keyed by the *external* input port identifier (the one + * the user-visible LogicalOp exposes). Each path points to a + * `.jsonl` file with a sibling `.schema.json`. + * @param outputDir destination directory; created if missing. Output files + * are named `output_port_<id>.jsonl` per external output. + */ + def execute( + opDesc: LogicalOp, + inputs: Map[PortIdentity, Path], + outputDir: Path + ): Result = { + Files.createDirectories(outputDir) + + // 1. Compile OpDesc → PhysicalPlan. For most ops this is a single-PhysicalOp + // plan; HashJoin and other multi-stage ops return multiple PhysicalOps + // plus internal PhysicalLinks (e.g. build.out → probe.in0). + val plan = opDesc.getPhysicalPlan(TestWorkflowId, TestExecutionId) + + // 2. Identify external input ports. A PhysicalOp port is "external" iff no + // PhysicalLink in this plan terminates at it. The user's `inputs` map + // must cover exactly these (matched by PortIdentity). + val externalInputs: Set[(PhysicalOpIdentity, PortIdentity)] = + plan.operators.flatMap { phOp => + phOp.inputPorts.keys.collect { + case portId if !plan.links.exists(l => l.toOpId == phOp.id && l.toPortId == portId) => + (phOp.id, portId) + } + } + validateInputCoverage(externalInputs, inputs.keySet) + + // 3. Load each input file as (Schema, Iterator[Tuple]). We read schemas + // eagerly but tuples lazily (saves memory on large fixtures). + val inputSchemas: Map[PortIdentity, Schema] = + inputs.map { case (portId, path) => portId -> TupleIO.readSchemaSidecar(path) } + val inputTuples: Map[PortIdentity, () => Iterator[Tuple]] = + inputs.map { + case (portId, path) => + val schema = inputSchemas(portId) + portId -> (() => TupleIO.readTuples(path, schema)) + } + + // 4. Propagate schemas. We CAN'T use `plan.propagateSchema(inputSchemas)` + // directly because it indexes by PortIdentity globally — for HashJoin + // probe's internal port 0 would collide with build's external port 0. + // Instead, set schemas only on the external (phOpId, portId) pairs and + // let `addLink` propagate the internal links' schemas naturally. + val planWithSchemas = + propagateExternalSchemas(plan, externalInputs, inputSchemas) + + // 5. Identify external output ports symmetrically: no outgoing PhysicalLink. + val externalOutputs: Set[(PhysicalOpIdentity, PortIdentity)] = + planWithSchemas.operators.flatMap { phOp => + phOp.outputPorts.keys.collect { + case portId + if !planWithSchemas.links + .exists(l => l.fromOpId == phOp.id && l.fromPortId == portId) => + (phOp.id, portId) + } + } + + // 6. Instantiate OpExec per PhysicalOp. We only support OpExecWithClassName; + // fail loudly otherwise so test authors know to mock Python UDFs out. + val opExecs: Map[PhysicalOpIdentity, OperatorExecutor] = + planWithSchemas.operators.map { phOp => + phOp.opExecInitInfo match { + case OpExecWithClassName(className, descString) => + phOp.id -> ExecFactory.newExecFromJavaClassName( + className, + descString, + idx = 0, + workerCount = 1 + ) + case other => + throw new UnsupportedOperationException( + s"OpExecHarness only supports OpExecWithClassName, got: $other" + ) + } + }.toMap + + // 7. Drive each PhysicalOp in topological order. Buffer outputs in memory + // keyed by (producer phOpId, output port). Downstream PhysicalOps then + // consume from these buffers via the plan's internal links. + val producedBuffer = + mutable.Map.empty[(PhysicalOpIdentity, PortIdentity), mutable.ArrayBuffer[Tuple]] + + planWithSchemas.topologicalIterator().foreach { phOpId => + val phOp = planWithSchemas.getOperator(phOpId) + val opExec = opExecs(phOpId) + runOneOp( + phOp, + opExec, + externalInputProvider = portId => inputTuples.get(portId).map(_.apply()), + upstreamBuffer = producedBuffer, + plan = planWithSchemas, + produced = producedBuffer + ) + } + + // 8. Materialize external outputs to JSONL with their propagated schemas. + val outputPaths = mutable.Map.empty[PortIdentity, Path] + val outputSchemas = mutable.Map.empty[PortIdentity, Schema] + externalOutputs.foreach { + case (phOpId, portId) => + val schema = planWithSchemas + .getOperator(phOpId) + .outputPorts(portId) + ._3 + .toOption + .getOrElse( + throw new IllegalStateException( + s"Output schema for ($phOpId, $portId) was not propagated" + ) + ) + val tuples = + producedBuffer.getOrElse((phOpId, portId), mutable.ArrayBuffer.empty[Tuple]) + val file = outputDir.resolve(s"output_port_${portId.id}.jsonl") + TupleIO.writeTuples(file, tuples.iterator, schema) + outputPaths(portId) = file + outputSchemas(portId) = schema + } + + Result(outputPaths.toMap, outputSchemas.toMap) + } + + /** + * Drives one PhysicalOp's lifecycle: open → input ports in dependency order + * (processTupleMultiPort + onFinishMultiPort per port) → close. Source ops + * (no input ports) get a single onFinishMultiPort(0) call which gives their + * `produceTuple()`-backed implementation a chance to emit. + * + * Outputs are bucketed by output PortIdentity. `processTupleMultiPort`'s + * `Option[PortIdentity]` return: None means port 0 (the default single- + * output convention used by the trait's fallback). Multi-output ops like + * Split set it explicitly. + */ + private def runOneOp( + phOp: PhysicalOp, + opExec: OperatorExecutor, + externalInputProvider: PortIdentity => Option[Iterator[Tuple]], + upstreamBuffer: mutable.Map[ + (PhysicalOpIdentity, PortIdentity), + mutable.ArrayBuffer[Tuple] + ], + plan: PhysicalPlan, + produced: mutable.Map[ + (PhysicalOpIdentity, PortIdentity), + mutable.ArrayBuffer[Tuple] + ] + ): Unit = { + opExec.open() + try { + def bucket(emitted: Iterator[(TupleLike, Option[PortIdentity])]): Unit = { + emitted.foreach { + case (tupleLike, portOpt) => + // Default: the op's single output port. Most operators have one + // output and use the trait's default port-0 wrapping, but + // multi-stage plans (e.g. HashJoin build) put their internal + // output on PortIdentity(0, internal = true) — a hardcoded + // PortIdentity(0, false) would NoSuchElementException here. + val outPortId = portOpt.getOrElse { + if (phOp.outputPorts.size == 1) phOp.outputPorts.keys.head + else PortIdentity(0) + } + val outSchema = phOp + .outputPorts(outPortId) + ._3 + .toOption + .getOrElse( + throw new IllegalStateException( + s"Op ${phOp.id} emitted to port $outPortId before its output schema was propagated" + ) + ) + val tuple = tupleLike + .asInstanceOf[org.apache.texera.amber.core.tuple.SeqTupleLike] + .enforceSchema(outSchema) + produced + .getOrElseUpdate((phOp.id, outPortId), mutable.ArrayBuffer.empty[Tuple]) += tuple + } + } + + // Process each input port in declared dependency order (e.g. HashJoin + // probe's build-side port must finish before the data-side port starts). + val portOrder = + if (phOp.getInputPortDependencyPairs.nonEmpty) + phOp.getInputPortDependencyPairs + else phOp.inputPorts.keys.toList.sortBy(_.id) + + portOrder.foreach { portId => + val tuples: Iterator[Tuple] = + if (externalInputProvider(portId).isDefined) { + externalInputProvider(portId).get + } else { + // Internal port: stitch upstream PhysicalLinks' buffers together + val upstream = plan.links + .filter(l => l.toOpId == phOp.id && l.toPortId == portId) + .toList + .sortBy(l => (l.fromOpId.toString, l.fromPortId.id)) + upstream.iterator + .flatMap(l => + upstreamBuffer + .getOrElse( + (l.fromOpId, l.fromPortId), + mutable.ArrayBuffer.empty[Tuple] + ) + .iterator + ) + } + + tuples.foreach { t => + bucket(opExec.processTupleMultiPort(t, portId.id)) + } + bucket(opExec.onFinishMultiPort(portId.id)) + } + + // Source operator: no input ports. Trigger production via onFinishMultiPort + // on a synthetic port 0 — SourceOperatorExecutor.onFinish ignores the port + // and emits everything from produceTuple(). + if (phOp.inputPorts.isEmpty) { + bucket(opExec.onFinishMultiPort(0)) + } + } finally { + opExec.close() + } + } + + // Walks the plan in topo order, propagating schemas only at the truly external + // input ports. Internal ports get their schema via `addLink` (PhysicalPlan + // re-applies the source's output schema to the destination port). This avoids + // a collision when multiple PhysicalOps share a PortIdentity (e.g. HashJoin + // probe.in0 internal vs build.in0 external both have PortIdentity(0)). + // `private[verify]` rather than private: the Python harness runs the same plan + // through a different executor, and how a plan's external ports get their + // schemas does not change with the executor behind them. + private[verify] def propagateExternalSchemas( + plan: PhysicalPlan, + externalPorts: Set[(PhysicalOpIdentity, PortIdentity)], + schemas: Map[PortIdentity, Schema] + ): PhysicalPlan = { + var acc = PhysicalPlan(operators = Set.empty, links = Set.empty) + plan.topologicalIterator().map(plan.getOperator).foreach { phOp => + val updated = phOp.inputPorts.keys.foldLeft(phOp) { (op, portId) => + if (externalPorts.contains((phOp.id, portId)) && schemas.contains(portId)) { + op.propagateSchema(Some((portId, schemas(portId)))) + } else op + } + // .propagateSchema() with no arg re-fires output derivation if all inputs + // are now resolved (source ops trigger immediately since inputPorts empty). + acc = acc.addOperator(updated.propagateSchema()) + plan.getUpstreamPhysicalLinks(phOp.id).foreach { link => + acc = acc.addLink(link) + } + } + acc + } + + private[verify] def validateInputCoverage( + external: Set[(PhysicalOpIdentity, PortIdentity)], + provided: Set[PortIdentity] + ): Unit = { + val expected = external.map(_._2) + val missing = expected -- provided + val extra = provided -- expected + require( + missing.isEmpty, + s"Missing input fixtures for external ports: $missing (expected $expected)" + ) + if (extra.nonEmpty) { + logger.warn(s"Input fixtures provided for non-external ports (ignored): $extra") + } + } +} + +/** + * JSON Lines I/O for Tuples. Each `.jsonl` file holds one record per line and + * is paired with a `.jsonl.schema.json` sidecar listing [[Attribute]]s in + * column order, which is what carries the types a JSON line cannot. + * + * pandas symmetry: `pd.read_json(path, lines=True)` and + * `df.to_json(path, orient='records', lines=True)` round-trip cleanly for the + * supported types (STRING / INTEGER / LONG / DOUBLE / BOOLEAN). + */ +object TupleIO { + + private def sidecar(path: Path): Path = + path.resolveSibling(path.getFileName.toString + ".schema.json") + + def readSchemaSidecar(path: Path): Schema = { + val text = new String(Files.readAllBytes(sidecar(path))) + objectMapper.readValue(text, classOf[Schema]) + } + + def readTuples(path: Path, schema: Schema): Iterator[Tuple] = { + // readAllLines closes the underlying handle; safer than Files.lines for + // test-scale fixtures where memory cost is negligible. + val lines = Files.readAllLines(path).asScala + lines.iterator.filter(_.trim.nonEmpty).map { line => + val node = objectMapper.readTree(line) + val builder = Tuple.builder(schema) + schema.getAttributes.foreach { attr => + val fieldNode = node.get(attr.getName) + val v: Any = + if (fieldNode == null || fieldNode.isNull) null + else + attr.getType match { + case AttributeType.STRING => fieldNode.asText() + case AttributeType.INTEGER => Int.box(fieldNode.asInt()) + case AttributeType.LONG => Long.box(fieldNode.asLong()) + case AttributeType.DOUBLE => Double.box(fieldNode.asDouble()) + case AttributeType.BOOLEAN => Boolean.box(fieldNode.asBoolean()) + case AttributeType.BINARY => + Base64.getDecoder.decode(fieldNode.asText()) + // Timestamps round-trip through the JDBC string form + // ("yyyy-mm-dd hh:mm:ss[.f]"), the exact inverse of Timestamp.toString + // below — timezone-free, so no shift across write/read. The Python + // side reads this column with convert_dates=False (see + // StandaloneRunner) and treats it as an opaque string, so both paths + // agree on pass-through. + case AttributeType.TIMESTAMP => + Timestamp.valueOf(fieldNode.asText()) + case other => + throw new UnsupportedOperationException( + s"TupleIO MVP doesn't support $other yet" + ) + } + builder.add(attr, v) + } + builder.build() + } + } + + def writeTuples(path: Path, tuples: Iterator[Tuple], schema: Schema): Unit = { + // Sidecar first so a partial main-file write still has a recoverable schema. + Files.write(sidecar(path), objectMapper.writeValueAsBytes(schema)) + val writer = Files.newBufferedWriter(path) + try { + tuples.foreach { t => + val node: ObjectNode = objectMapper.createObjectNode() + schema.getAttributes.zipWithIndex.foreach { + case (attr, idx) => + val v = t.getField[Any](idx) + if (v == null) node.putNull(attr.getName) + else + attr.getType match { + case AttributeType.STRING => node.put(attr.getName, v.toString) + case AttributeType.INTEGER => node.put(attr.getName, v.asInstanceOf[Int]) + case AttributeType.LONG => node.put(attr.getName, v.asInstanceOf[Long]) + case AttributeType.DOUBLE => node.put(attr.getName, v.asInstanceOf[Double]) + case AttributeType.BOOLEAN => node.put(attr.getName, v.asInstanceOf[Boolean]) + case AttributeType.BINARY => + node.put( + attr.getName, + Base64.getEncoder.encodeToString(v.asInstanceOf[Array[Byte]]) + ) + case AttributeType.TIMESTAMP => + node.put(attr.getName, v.asInstanceOf[Timestamp].toString) + case other => + throw new UnsupportedOperationException( + s"TupleIO MVP doesn't support $other yet" + ) + } + } + writer.write(objectMapper.writeValueAsString(node)) + writer.newLine() + } + } finally writer.close() + } +}
