Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
@@ -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
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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)")
Expand Down Expand Up @@ -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"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
}
}
Loading
Loading