pub struct ReduceSpec { /* private fields */ }Expand description
A reduce operation over already-sorted input.
Every input table must already be sorted by a column set that starts with
reduce_by — SortSpec is how a table gets that way. When it is,
this is the operation to reach for: a map-reduce over the same data would
pay for a shuffle that has already been done.
use ytsaurus_client::ReduceSpec;
let spec = ReduceSpec::new("./wordcount reduce", ["//tmp/sorted"], ["//tmp/counts"], ["word"])
.with_local_file("//tmp/wordcount");Implementations§
Source§impl ReduceSpec
impl ReduceSpec
Sourcepub fn new<I, O, K>(
command: impl Into<String>,
inputs: I,
outputs: O,
reduce_by: K,
) -> Selfwhere
I: IntoIterator,
I::Item: Into<String>,
O: IntoIterator,
O::Item: Into<String>,
K: IntoIterator,
K::Item: Into<String>,
pub fn new<I, O, K>(
command: impl Into<String>,
inputs: I,
outputs: O,
reduce_by: K,
) -> Selfwhere
I: IntoIterator,
I::Item: Into<String>,
O: IntoIterator,
O::Item: Into<String>,
K: IntoIterator,
K::Item: Into<String>,
A reduce running command over inputs, grouped by reduce_by.
Sourcepub fn with_local_file(self, path: impl Into<String>) -> Self
pub fn with_local_file(self, path: impl Into<String>) -> Self
Adds a Cypress file the job needs — normally the worker binary.
Sourcepub fn with_local_file_named(
self,
path: impl Into<String>,
name: impl AsRef<str>,
) -> Self
pub fn with_local_file_named( self, path: impl Into<String>, name: impl AsRef<str>, ) -> Self
Adds a Cypress file under a different name in the job’s sandbox.
See MapSpec::with_local_file_named for why the name matters.
Sourcepub fn with_memory_limit(self, bytes: i64) -> Self
pub fn with_memory_limit(self, bytes: i64) -> Self
Sets the reducer’s memory limit, in bytes.
Sourcepub fn with_formats(self, input: DataFormat, output: DataFormat) -> Self
pub fn with_formats(self, input: DataFormat, output: DataFormat) -> Self
Selects the reducer’s input and output data formats.
YSON selections apply to every table. A Skiff selection must contain one table schema per corresponding input or output table, in the same order. The default remains binary YSON.
A Skiff reducer receives its key switch as a $key_switch boolean
column rather than as a YSON control record, so the input schema has to
declare that column for ytsaurus-job’s SkiffJobReader to report it —
enable_key_switch asks the cluster to deliver key switches, and the
format decides how they arrive. A schema without the column leaves a
grouping reducer seeing one group, exactly as
Self::without_key_switch would.
Sourcepub fn with_skiff_formats(self, input: SkiffFormat, output: SkiffFormat) -> Self
pub fn with_skiff_formats(self, input: SkiffFormat, output: SkiffFormat) -> Self
Uses validated Skiff formats for the reducer’s input and output streams.
This compatibility convenience delegates to Self::with_formats.
Sourcepub fn with_env(self, key: impl Into<String>, value: impl Into<String>) -> Self
pub fn with_env(self, key: impl Into<String>, value: impl Into<String>) -> Self
Sets an environment variable for the job, e.g. RUST_BACKTRACE.
Sourcepub fn with_sort_by<K>(self, columns: K) -> Self
pub fn with_sort_by<K>(self, columns: K) -> Self
Sets the columns the input is sorted by, when they differ from
reduce_by.
reduce_by must be a prefix of them. Saying so asks the cluster to
check the input really is sorted that way, and guarantees the order rows
arrive in within a group.
Sourcepub fn with_job_count(self, count: i64) -> Self
pub fn with_job_count(self, count: i64) -> Self
Requests a specific job count.
Sourcepub fn with_input_table_index(self) -> Self
pub fn with_input_table_index(self) -> Self
Asks for the input table index to be delivered with each row.
Reduce merges several sorted tables into one stream, so this is how a job tells which table a row came from.
Sourcepub fn without_key_switch(self) -> Self
pub fn without_key_switch(self) -> Self
Turns off key_switch delivery to the reducer.
Sourcepub fn skiff_table_mismatch(&self) -> Option<String>
pub fn skiff_table_mismatch(&self) -> Option<String>
Describes a Skiff format that does not match this spec’s table lists.
A reduce merges its input tables into one sorted stream but keeps them
distinguishable, so the input format describes every input table, as the
Go SDK’s setupSkiffInputFormat also requires. See
MapSpec::skiff_table_mismatch for what an unchecked mismatch costs.
Client::start_reduce checks this before
sending the spec.
Trait Implementations§
Source§impl Clone for ReduceSpec
impl Clone for ReduceSpec
Source§fn clone(&self) -> ReduceSpec
fn clone(&self) -> ReduceSpec
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read more