diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/SamplingHelpers.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/SamplingHelpers.scala new file mode 100644 index 00000000000..1fa9210898b --- /dev/null +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/SamplingHelpers.scala @@ -0,0 +1,68 @@ +/* + * 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.operator + +/** The Python a sampler's standalone code needs, emitted once per script via + * [[StandaloneCodeGenerator.standaloneHelpers]]. + */ +object SamplingHelpers { + + /** + * A Python transcription of `java.util.Random`, for operators whose executor + * draws from one. + * + * A sampler decides per row whether to keep it, so which rows survive is + * fixed by the exact sequence the generator produces. Seeding Python's + * `random` or numpy's with the engine's seed selects a different set, and + * the script would then report a different sample than the workflow it came + * from. Only the same generator gives the same rows. + */ + val JavaRandom: String = + """# java.util.Random, transcribed so sampling matches the engine. + |class _TexeraJavaRandom: + | _MASK = (1 << 48) - 1 + | _MULTIPLIER = 0x5DEECE66D + | _ADDEND = 0xB + | + | def __init__(self, seed): + | self._seed = (seed ^ self._MULTIPLIER) & self._MASK + | + | def _next(self, bits): + | self._seed = (self._seed * self._MULTIPLIER + self._ADDEND) & self._MASK + | value = self._seed >> (48 - bits) + | return value - (1 << 32) if value >= (1 << 31) else value + | + | def next_double(self): + | return ((self._next(26) << 27) + self._next(27)) * (2.0 ** -53) + | + | def next_int(self, bound): + | if bound <= 0: + | raise ValueError("bound must be positive") + | if bound & (-bound) == bound: + | return (bound * self._next(31)) >> 31 + | while True: + | bits = self._next(31) + | value = bits % bound + | # Java rejects a draw by letting this sum overflow an int. + | # Python would carry it and never reject. + | probe = bits - value + (bound - 1) + | if ((probe + (1 << 31)) % (1 << 32)) - (1 << 31) >= 0: + | return value""".stripMargin +} diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/dummy/DummyOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/dummy/DummyOpDesc.scala index 4a2c4c48d54..bee942cd370 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/dummy/DummyOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/dummy/DummyOpDesc.scala @@ -23,9 +23,14 @@ import com.fasterxml.jackson.annotation.{JsonProperty, JsonPropertyDescription} import com.kjetland.jackson.jsonSchema.annotations.JsonSchemaTitle import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PortIdentity} import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} -import org.apache.texera.amber.operator.{LogicalOp, PortDescription, PortDescriptor} +import org.apache.texera.amber.operator.{ + LogicalOp, + PortDescription, + PortDescriptor, + StandaloneCodeGenerator +} -class DummyOpDesc extends LogicalOp with PortDescriptor { +class DummyOpDesc extends LogicalOp with PortDescriptor with StandaloneCodeGenerator { @JsonProperty @JsonSchemaTitle("Description") @@ -66,4 +71,12 @@ class DummyOpDesc extends LogicalOp with PortDescriptor { allowPortCustomization = true ) } + + override def generateStandaloneCode(): String = { + // Placeholder operator: pass the first input through to the output. + // Multi-port configurations don't fully translate under the current + // single-output placeholder scheme; downstream of extra ports would + // alias the same variable. + "out1df = in1df" + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/ifStatement/IfOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/ifStatement/IfOpDesc.scala index f811ed34296..296e70e6928 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/ifStatement/IfOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/ifStatement/IfOpDesc.scala @@ -24,11 +24,12 @@ import com.kjetland.jackson.jsonSchema.annotations.JsonSchemaTitle import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow._ -import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} +import org.apache.texera.amber.pybuilder.PythonTemplateBuilder.pyStringLiteral import org.apache.texera.amber.util.JSONUtils.objectMapper -class IfOpDesc extends LogicalOp { +class IfOpDesc extends LogicalOp with StandaloneCodeGenerator { @JsonProperty(required = true) @JsonSchemaTitle("Condition State") @JsonPropertyDescription("name of the state variable to evaluate") @@ -72,4 +73,33 @@ class IfOpDesc extends LogicalOp { ), outputPorts = List(OutputPort(PortIdentity(), "False"), OutputPort(PortIdentity(1), "True")) ) + + // The engine flips the active output port when a State message arrives on the + // Condition port. A script has no State channel, so the condition is read from + // a global named after it, true when nothing set it, which is the port the + // engine starts on. The Condition input carries no rows to read either way. + // + // Taking True silently would report the True route for a workflow whose + // condition says otherwise, so an unset switch says which route it took and + // how to take the other. + override def standaloneImports(): Seq[String] = Seq("import sys") + + override def generateStandaloneCode(): String = { + val globalName = "_texera_if_" + Option(conditionName).getOrElse("") + val globalLit = pyStringLiteral(globalName) + s"""if $globalLit not in globals(): + | print( + | "If: the script cannot receive the condition the workflow sets upstream, " + | "so it takes the True branch. To take the False branch, set " + | + $globalLit + " = False before this step.", + | file=sys.stderr, + | ) + |_texera_if_cond = bool(globals().get($globalLit, True)) + |if _texera_if_cond: + | out2df = in2df.copy() + | out1df = in2df.iloc[0:0].copy() + |else: + | out1df = in2df.copy() + | out2df = in2df.iloc[0:0].copy()""".stripMargin + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDesc.scala index 31cafcbc76b..188d3e38ef4 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDesc.scala @@ -23,11 +23,12 @@ import com.fasterxml.jackson.annotation.{JsonProperty, JsonPropertyDescription} import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PhysicalOp} +import org.apache.texera.amber.operator.{SamplingHelpers, StandaloneCodeGenerator} import org.apache.texera.amber.operator.filter.FilterOpDesc import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} import org.apache.texera.amber.util.JSONUtils.objectMapper -class RandomKSamplingOpDesc extends FilterOpDesc { +class RandomKSamplingOpDesc extends FilterOpDesc with StandaloneCodeGenerator { @JsonProperty(value = "random k sample percentage", required = true) @JsonPropertyDescription("random k sampling with given percentage") @@ -60,4 +61,23 @@ class RandomKSamplingOpDesc extends FilterOpDesc { outputPorts = List(OutputPort()), supportReconfiguration = true ) + + override def standaloneHelpers(): Seq[String] = Seq(SamplingHelpers.JavaRandom) + + // A per-row Bernoulli filter drawn from the same generator, so a rerun keeps + // the exact rows rather than merely the same proportion of them. + // + // The engine seeds it with the worker count, which makes the surviving rows a + // property of the deployment. One process states the one seed it can and so + // agrees with a single-worker run; a wider one splits the input and samples + // each share from the start of the sequence, which no seed here reproduces. + override def generateStandaloneCode(): String = { + val p = percentage / 100.0 + s"""_texera_rks_rng = _TexeraJavaRandom(1) + |_texera_rks_mask = pd.Series( + | [$p >= _texera_rks_rng.next_double() for _ in range(len(in1df))], + | index=in1df.index, dtype=bool + |) + |out1df = in1df[_texera_rks_mask].reset_index(drop=True)""".stripMargin + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDesc.scala index eda2aa299fe..a20112f84b6 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDesc.scala @@ -23,11 +23,11 @@ import com.fasterxml.jackson.annotation.{JsonProperty, JsonPropertyDescription} import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PhysicalOp} -import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.operator.{LogicalOp, SamplingHelpers, StandaloneCodeGenerator} import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} import org.apache.texera.amber.util.JSONUtils.objectMapper -class ReservoirSamplingOpDesc extends LogicalOp { +class ReservoirSamplingOpDesc extends LogicalOp with StandaloneCodeGenerator { @JsonProperty(value = "number of item sampled in reservoir sampling", required = true) @JsonPropertyDescription("reservoir sampling with k items being kept randomly") @@ -60,4 +60,32 @@ class ReservoirSamplingOpDesc extends LogicalOp { outputPorts = List(OutputPort()) ) } + + override def standaloneHelpers(): Seq[String] = Seq(SamplingHelpers.JavaRandom) + + // The executor runs Algorithm R (Vitter): fill the reservoir with the first k + // tuples, then for tuple m+1 (m >= k) draw i = rand.nextInt(m), uniform in + // [0, m), and replace reservoir[i] iff i < k. + // + // Its generator is seeded with the worker count, so the rows agree with a + // single-worker run and not with a wider one. See RandomKSamplingOpDesc for + // why no seed closes that gap. + // + // The reservoir holds row positions, and the rows are taken from the input + // by them, so every column keeps its dtype. A frame rebuilt from the rows' + // values infers each one again, and a timestamp past pandas' nanosecond + // range came back as an object column. + override def generateStandaloneCode(): String = { + s"""_texera_rs_rng = _TexeraJavaRandom(1) + |_texera_rs_k = $k + |_texera_rs_reservoir = [] + |for _texera_rs_n in range(len(in1df)): + | if _texera_rs_n < _texera_rs_k: + | _texera_rs_reservoir.append(_texera_rs_n) + | else: + | _texera_rs_i = _texera_rs_rng.next_int(_texera_rs_n) + | if _texera_rs_i < _texera_rs_k: + | _texera_rs_reservoir[_texera_rs_i] = _texera_rs_n + |out1df = in1df.iloc[_texera_rs_reservoir].reset_index(drop=True)""".stripMargin + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sleep/SleepOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sleep/SleepOpDesc.scala index 3eee3cd33b7..1986eb52d83 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sleep/SleepOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/sleep/SleepOpDesc.scala @@ -24,11 +24,11 @@ import com.kjetland.jackson.jsonSchema.annotations.JsonSchemaTitle import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PhysicalOp} -import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} import org.apache.texera.amber.util.JSONUtils.objectMapper -class SleepOpDesc extends LogicalOp { +class SleepOpDesc extends LogicalOp with StandaloneCodeGenerator { @JsonProperty(required = true) @JsonSchemaTitle("Sleep Time (seconds)") @@ -63,4 +63,10 @@ class SleepOpDesc extends LogicalOp { inputPorts = List(InputPort()), outputPorts = List(OutputPort()) ) + + override def generateStandaloneCode(): String = { + // JVM op sleeps between each tuple; pandas operates in batch, not row-by-row, + // so per-row sleep has no equivalent. Translate as a passthrough. + "out1df = in1df" + } } diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/split/SplitOpDesc.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/split/SplitOpDesc.scala index c3c34df948c..aa3c377da59 100644 --- a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/split/SplitOpDesc.scala +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/split/SplitOpDesc.scala @@ -29,12 +29,12 @@ import com.kjetland.jackson.jsonSchema.annotations.{ import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} import org.apache.texera.amber.core.workflow._ -import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.operator.{LogicalOp, SamplingHelpers, StandaloneCodeGenerator} import org.apache.texera.amber.operator.metadata.annotations.HideAnnotation import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} import org.apache.texera.amber.util.JSONUtils.objectMapper -class SplitOpDesc extends LogicalOp { +class SplitOpDesc extends LogicalOp with StandaloneCodeGenerator { @JsonSchemaTitle("Split Percentage") @JsonProperty(defaultValue = "80") @@ -99,4 +99,25 @@ class SplitOpDesc extends LogicalOp { ) } + override def standaloneHelpers(): Seq[String] = Seq(SamplingHelpers.JavaRandom) + + override def generateStandaloneCode(): String = { + // The executor sends each tuple to the upper port iff nextInt(100) < k, + // drawing from one generator per run. Reproducing that draw keeps the same + // rows on the same side; the mask drives both ports, so `out1df` takes the + // upper k% and `out2df` the remainder and the two partition the input. + // + // With "Auto-Generate Seed" the executor seeds from the clock, so that run + // is not reproducible by anything, itself included — the script then seeds + // from its own clock, matching the intent rather than a particular run. + val seedExpr = if (random) "int(_texera_time.time() * 1000)" else seed.toString + val timeImport = if (random) "import time as _texera_time\n" else "" + s"""${timeImport}_texera_split_rng = _TexeraJavaRandom($seedExpr) + |_texera_split_mask = pd.Series( + | [_texera_split_rng.next_int(100) < $k for _ in range(len(in1df))], + | index=in1df.index, dtype=bool + |) + |out1df = in1df[_texera_split_mask].reset_index(drop=True) + |out2df = in1df[~_texera_split_mask].reset_index(drop=True)""".stripMargin + } } diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/SamplingHelpersSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/SamplingHelpersSpec.scala new file mode 100644 index 00000000000..8bf4f8f6f22 --- /dev/null +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/SamplingHelpersSpec.scala @@ -0,0 +1,117 @@ +/* + * 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.operator + +import com.typesafe.config.ConfigFactory +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 + +/** The transcribed java.util.Random, checked against what Java would answer. + * + * Only the generator is exercised here. What each sampling operator does with + * the numbers is its own spec's business. + */ +class SamplingHelpersSpec extends AnyFlatSpec with Matchers { + + /** Run `body` after the helper class, and return what it printed. */ + private def runWithHelper(body: String): (Int, String) = { + val python = resolvePython().getOrElse(cancel("No runnable python executable")) + val script = Files.createTempFile("sampling-helpers-", ".py") + script.toFile.deleteOnExit() + Files.write( + script, + (SamplingHelpers.JavaRandom + "\n\n" + body).getBytes(StandardCharsets.UTF_8) + ) + val process = new ProcessBuilder(python, script.toString).redirectErrorStream(true).start() + val out = Source.fromInputStream(process.getInputStream).mkString + process.waitFor(60, TimeUnit.SECONDS) + (process.exitValue(), out) + } + + // A bound of 2**30 + 1 puts about half the draws in the rejection zone, so the + // first one with this seed is already a draw Java throws away. The five + // numbers are what java.util.Random answers. + "SamplingHelpers.JavaRandom" should "reject the draws Java's overflow check rejects" in { + val (exit, out) = runWithHelper( + """r = _TexeraJavaRandom(1) + |print([r.next_int((1 << 30) + 1) for _ in range(5)]) + |""".stripMargin + ) + withClue(s"python said:\n$out") { + exit shouldBe 0 + out.trim shouldBe "[215764588, 880641847, 874970313, 446064254, 77814904]" + } + } + + // Reservoir Sampling reaches this check with a reservoir of zero, and the + // engine ends the run there. + it should "refuse a bound that is not positive, the way Java does" in { + val (exit, out) = runWithHelper( + """r = _TexeraJavaRandom(1) + |try: + | r.next_int(0) + | print("no error") + |except ValueError as e: + | print("ValueError:", e) + |""".stripMargin + ) + withClue(s"python said:\n$out") { + exit shouldBe 0 + out.trim shouldBe "ValueError: bound must be positive" + } + } + + // The power-of-two bound takes the branch that never rejects. + it should "take the power-of-two shortcut without rejecting" in { + val (exit, out) = runWithHelper( + """r = _TexeraJavaRandom(42) + |print(all(0 <= r.next_int(256) < 256 for _ in range(1000))) + |""".stripMargin + ) + withClue(s"python said:\n$out") { + exit shouldBe 0 + out.trim shouldBe "True" + } + } + + 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) + } +} diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/dummy/DummyOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/dummy/DummyOpDescSpec.scala index 6f23261b271..4ee8cfe7fce 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/dummy/DummyOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/dummy/DummyOpDescSpec.scala @@ -71,6 +71,11 @@ class DummyOpDescSpec extends AnyFlatSpec with Matchers { (new DummyOpDesc).dummyOperator shouldBe "" } + "DummyOpDesc.generateStandaloneCode" should + "pass the first input through to the output" in { + (new DummyOpDesc).generateStandaloneCode() shouldBe "out1df = in1df" + } + "DummyOpDesc.getPhysicalOp" should "be the unimplemented LogicalOp stub (throws NotImplementedError)" in { intercept[NotImplementedError] { diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/ifStatement/IfOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/ifStatement/IfOpDescSpec.scala index 49cd325d572..03e879caf06 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/ifStatement/IfOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/ifStatement/IfOpDescSpec.scala @@ -19,16 +19,27 @@ package org.apache.texera.amber.operator.ifStatement +import com.typesafe.config.ConfigFactory import org.apache.texera.amber.core.executor.OpExecWithClassName 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.PortIdentity import org.apache.texera.amber.operator.metadata.OperatorGroupConstants +import org.apache.texera.amber.operator.tags.IntegrationTest +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 IfOpDescSpec extends AnyFlatSpec with Matchers { + private val NeedsPythonPackages = Tag(classOf[IntegrationTest].getName) + private val workflowId = WorkflowIdentity(1L) private val executionId = ExecutionIdentity(1L) @@ -76,4 +87,82 @@ class IfOpDescSpec extends AnyFlatSpec with Matchers { ) out shouldBe Map(PortIdentity() -> dataSchema, PortIdentity(1) -> dataSchema) } + + "IfOpDesc.generateStandaloneCode" should + "name the switch after the condition so two Ifs in one script do not share it" in { + val op = new IfOpDesc + op.conditionName = "ready" + op.generateStandaloneCode() should include("\"_texera_if_ready\"") + } + + // The engine picks the route from a State message, which the verification + // harness has no channel for, so False is only reachable here. + it should "send the rows to True by default and to False when the switch is off" taggedAs NeedsPythonPackages in { + val python = resolvePythonExecutable().getOrElse(cancel("No runnable python executable")) + if (!canImportPandas(python)) cancel(s"'$python' cannot import pandas") + + val op = new IfOpDesc + op.conditionName = "ready" + val body = op.generateStandaloneCode().linesIterator.map(" " + _).mkString("\n") + + val driver = + s"""import pandas as pd + |import sys + | + |for _switch in (None, True, False): + | in2df = pd.DataFrame({"id": [1, 2, 3]}) + | if _switch is None: + | globals().pop("_texera_if_ready", None) + | else: + | globals()["_texera_if_ready"] = _switch + |$body + | print(_switch, list(out1df["id"]), list(out2df["id"])) + |""".stripMargin + + val script = Files.createTempFile("if-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 + val lines = out.trim.linesIterator.toSeq + // out1df is the False port, out2df the True port. Each route has to be + // exclusive: the rows leave by one and the other comes back empty. + lines should contain("None [] [1, 2, 3]") + lines should contain("True [] [1, 2, 3]") + lines should contain("False [1, 2, 3] []") + // Only the unset switch takes True without being told to, so it alone + // says so. + lines.count(_.contains("_texera_if_ready = False before this step")) shouldBe 1 + } + } + + private def resolvePythonExecutable(): 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 isRunnable(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(isRunnable) + } + + 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/randomksampling/RandomKSamplingOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDescSpec.scala index 94d450a9559..7d6f3781abb 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/randomksampling/RandomKSamplingOpDescSpec.scala @@ -59,6 +59,30 @@ class RandomKSamplingOpDescSpec extends AnyFlatSpec with Matchers { restored.asInstanceOf[RandomKSamplingOpDesc].percentage shouldBe 25 } + "RandomKSamplingOpDesc.generateStandaloneCode" should + "draw a per-row Bernoulli mask from a seeded RNG" in { + val d = new RandomKSamplingOpDesc + d.percentage = 25 + d.generateStandaloneCode() shouldBe + """_texera_rks_rng = _TexeraJavaRandom(1) + |_texera_rks_mask = pd.Series( + | [0.25 >= _texera_rks_rng.next_double() for _ in range(len(in1df))], + | index=in1df.index, dtype=bool + |) + |out1df = in1df[_texera_rks_mask].reset_index(drop=True)""".stripMargin + } + + // percentage is a whole number on the descriptor but a probability in the + // generated code; an integer division here would keep no rows at all. + it should "convert the percentage to a probability" in { + Seq(0 -> "0.0", 30 -> "0.3", 100 -> "1.0").foreach { + case (pct, prob) => + val d = new RandomKSamplingOpDesc + d.percentage = pct + d.generateStandaloneCode() should include(s"[$prob >= _texera_rks_rng.next_double()") + } + } + "RandomKSamplingOpDesc.getPhysicalOp" should "wire the RandomKSamplingOpExec class name and carry ports" in { val d = new RandomKSamplingOpDesc diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDescSpec.scala index 0c903d6ce88..0e806034f64 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/reservoirsampling/ReservoirSamplingOpDescSpec.scala @@ -19,16 +19,27 @@ package org.apache.texera.amber.operator.reservoirsampling +import com.typesafe.config.ConfigFactory import org.apache.texera.amber.core.executor.OpExecWithClassName import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} -import org.apache.texera.amber.operator.LogicalOp +import org.apache.texera.amber.operator.tags.IntegrationTest +import org.apache.texera.amber.operator.{LogicalOp, SamplingHelpers} import org.apache.texera.amber.operator.metadata.OperatorGroupConstants 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 ReservoirSamplingOpDescSpec extends AnyFlatSpec with Matchers { + private val NeedsPythonPackages = Tag(classOf[IntegrationTest].getName) + private val workflowId = WorkflowIdentity(1L) private val executionId = ExecutionIdentity(1L) @@ -59,6 +70,56 @@ class ReservoirSamplingOpDescSpec extends AnyFlatSpec with Matchers { restored.asInstanceOf[ReservoirSamplingOpDesc].k shouldBe 100 } + // Algorithm R: fill the reservoir with the first k rows, then replace a + // uniformly-drawn slot for each later row. Pin the whole snippet — the + // generated Python is indentation-sensitive. + "ReservoirSamplingOpDesc.generateStandaloneCode" should + "emit a seeded Algorithm R reservoir over the input rows" in { + val d = new ReservoirSamplingOpDesc + d.k = 3 + d.generateStandaloneCode() shouldBe + """_texera_rs_rng = _TexeraJavaRandom(1) + |_texera_rs_k = 3 + |_texera_rs_reservoir = [] + |for _texera_rs_n in range(len(in1df)): + | if _texera_rs_n < _texera_rs_k: + | _texera_rs_reservoir.append(_texera_rs_n) + | else: + | _texera_rs_i = _texera_rs_rng.next_int(_texera_rs_n) + | if _texera_rs_i < _texera_rs_k: + | _texera_rs_reservoir[_texera_rs_i] = _texera_rs_n + |out1df = in1df.iloc[_texera_rs_reservoir].reset_index(drop=True)""".stripMargin + } + + // A reservoir of zero ends the engine's run on the first row: the executor + // skips the fill branch and hands nextInt a bound of zero, which Java refuses. + // The script has to refuse it there too. Answering with an empty table would + // report a result the run never produced. + it should "fail on the first row when the reservoir holds nothing" taggedAs NeedsPythonPackages in { + val python = resolvePython().getOrElse(cancel("No runnable python executable")) + if (!canImportPandas(python)) cancel(s"'$python' cannot import pandas") + + val d = new ReservoirSamplingOpDesc + d.k = 0 + val script = Files.createTempFile("reservoir-zero-", ".py") + script.toFile.deleteOnExit() + val driver = + s"""import pandas as pd + |${SamplingHelpers.JavaRandom} + |in1df = pd.DataFrame({"id": [1, 2, 3]}) + |${d.generateStandaloneCode()} + |""".stripMargin + 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(60, TimeUnit.SECONDS) + withClue(s"python said:\n$out\nscript:\n$driver") { + process.exitValue() should not be 0 + out should include("bound must be positive") + } + } + "ReservoirSamplingOpDesc.getPhysicalOp" should "wire the ReservoirSamplingOpExec class name and carry ports" in { val d = new ReservoirSamplingOpDesc @@ -73,4 +134,30 @@ class ReservoirSamplingOpDescSpec extends AnyFlatSpec with Matchers { physical.inputPorts.keySet shouldBe d.operatorInfo.inputPorts.map(_.id).toSet physical.outputPorts.keySet shouldBe d.operatorInfo.outputPorts.map(_.id).toSet } + + 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/sleep/SleepOpExecSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sleep/SleepOpExecSpec.scala index 2aa143f2948..c5280d30994 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sleep/SleepOpExecSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/sleep/SleepOpExecSpec.scala @@ -54,4 +54,21 @@ class SleepOpExecSpec extends AnyFlatSpec { val emitted = (0 until 5).flatMap(i => exec.processTuple(tuple(i), 0).toList) assert(emitted == (0 until 5).map(tuple).toList) } + + // The one thing the operator is for. Every other case here configures no delay, + // so the seconds-to-millis conversion was the only line nothing asked about: + // dropping it, or losing the factor of a thousand, left them all passing. + // + // One second is the smallest delay the field can ask for, sleepTime being whole + // seconds, and one tuple is enough to observe it. The lower bound is asserted + // rather than a window: a loaded machine may take longer, and a delay that + // overshoots is not the failure this guards against. + it should "delay each tuple by the configured number of seconds" in { + val exec = new SleepOpExec(descString(1)) + val startNanos = System.nanoTime() + val emitted = exec.processTuple(tuple(1), 0).toList + val elapsedMillis = (System.nanoTime() - startNanos) / 1000000 + assert(emitted == List(tuple(1))) + assert(elapsedMillis >= 1000, s"expected at least 1000 ms of sleep, took $elapsedMillis ms") + } } diff --git a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/split/SplitOpDescSpec.scala b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/split/SplitOpDescSpec.scala index d1fecdf7834..d0a769e9472 100644 --- a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/split/SplitOpDescSpec.scala +++ b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/split/SplitOpDescSpec.scala @@ -126,6 +126,44 @@ class SplitOpDescSpec extends AnyFlatSpec with Matchers { } } + // --------------------------------------------------------------------------- + // generateStandaloneCode + // --------------------------------------------------------------------------- + + // The generator is a transcription of java.util.Random, so the script draws the + // rows the executor drew rather than a differently-shaped sample from the same + // distribution. With auto-generate on, the seed comes from the clock, as it does + // on the executor. + "SplitOpDesc.generateStandaloneCode" should + "draw an unseeded mask and partition the input across both outputs" in { + (new SplitOpDesc).generateStandaloneCode() shouldBe + """import time as _texera_time + |_texera_split_rng = _TexeraJavaRandom(int(_texera_time.time() * 1000)) + |_texera_split_mask = pd.Series( + | [_texera_split_rng.next_int(100) < 80 for _ in range(len(in1df))], + | index=in1df.index, dtype=bool + |) + |out1df = in1df[_texera_split_mask].reset_index(drop=True) + |out2df = in1df[~_texera_split_mask].reset_index(drop=True)""".stripMargin + } + + it should "pass the configured seed to the generator when auto-generate is off" in { + val d = new SplitOpDesc + d.random = false + d.seed = 42 + d.generateStandaloneCode() should include("_texera_split_rng = _TexeraJavaRandom(42)") + } + + // k is a percentage on both sides, compared as the executor compares it: + // nextInt(100) < k, not a uniform draw against k/100. + it should "compare the draw against the percentage the executor uses" in { + Seq(0, 30, 80, 100).foreach { k => + val d = new SplitOpDesc + d.k = k + d.generateStandaloneCode() should include(s"_texera_split_rng.next_int(100) < $k") + } + } + // --------------------------------------------------------------------------- // Independent instances // ---------------------------------------------------------------------------