plan API

包 Luna-Flow/luna_thread/plan 以 @plan 导入,把一个数据并行操作描述为一个值:它的种类、输入的元素类型和长度、执行策略、可选的归约内核以及顺序保证。它按 v1 子集验证计划。它不执行任何东西,可以在所有目标上构建。

本包的枚举类型在包外是只读的:可以对其构造器做模式匹配,但请用下面的函数构建值。

类型

PlanKind

计划可以描述的四种操作。

pub enum PlanKind {
  Map
  Reduce
  Scan
  MapReduce
} derive(Eq, @debug.Debug)

Map 对每个元素应用一个函数,Reduce 用内核合并所有元素,Scan 产生逐步累积的结果(前缀),MapReduce 先映射再归约。

ValueType

计划输入的元素类型。

pub enum ValueType {
  I32
  I64
  F32
  F64
  Bytes
  Opaque(String)
} derive(Eq, @debug.Debug)

v1 只支持 I32 和 I64。

ReductionKernel

归约、映射归约或扫描所用的二元运算。

pub enum ReductionKernel {
  Sum
  Min
  Max
  Custom(String)
} derive(Eq, @debug.Debug)

Custom(name) 指一个运行时不认识的内核;在 v1 中它永远无法通过验证。

DataDomain

计划输入的元素类型和长度。

pub struct DataDomain {
  element_type : ValueType
  input_length : Int
} derive(Eq, @debug.Debug)
pub fn DataDomain::new(ValueType, Int) -> Self
pub fn DataDomain::element_type(Self) -> ValueType
pub fn DataDomain::input_length(Self) -> Int

DataDomain::new(t, n) 不检查地保存参数;validate 会拒绝 n≤0n \le 0。

Plan

一个完整的计划。

pub struct Plan {
  label : String
  kind : PlanKind
  domain : DataDomain
  policy : @shared.ExecutionPolicy
  reduction : ReductionKernel?
  ordering : @shared.OrderingGuarantee
} derive(Eq, @debug.Debug)
pub fn Plan::new(String, PlanKind, DataDomain, policy? : @shared.ExecutionPolicy, reduction? : ReductionKernel, ordering? : @shared.OrderingGuarantee) -> Self

Plan::new 不检查地保存参数。policy 默认为 @shared.native_policy(),reduction 默认为 None,ordering 默认为 PreserveInputOrder。下面的构建函数无法构造的组合(例如顺序宽松的扫描)可以用它来构造。

Plan::label、Plan::kind、Plan::domain、Plan::policy、Plan::reduction 和 Plan::ordering

返回计划的各个字段。

pub fn Plan::label(Self) -> String
pub fn Plan::kind(Self) -> PlanKind
pub fn Plan::domain(Self) -> DataDomain
pub fn Plan::policy(Self) -> @shared.ExecutionPolicy
pub fn Plan::reduction(Self) -> ReductionKernel?
pub fn Plan::ordering(Self) -> @shared.OrderingGuarantee

ValidationIssue

计划超出 v1 子集的一个原因。

pub enum ValidationIssue {
  EmptyInput
  ChunkSizeDoesNotFitInput(chunk_size~ : Int, input_length~ : Int)
  WorkerCountExceedsInput(worker_count~ : Int, input_length~ : Int)
  MissingReductionKernel
  ScanRequiresStableOrdering
  UnsupportedValueType
  UnsupportedReductionKernel
  UnsupportedBackend
  UnsupportedMode
  UnsupportedOrdering(kind~ : PlanKind, ordering~ : @shared.OrderingGuarantee)
} derive(Eq, @debug.Debug)

validate 说明了每种问题在何时报告。

T::equal

结构相等,作为方法提升到本包的每个类型上:DataDomain、Plan、PlanKind、ReductionKernel、ValidationIssue 和 ValueType。请使用 == 和 !=。

pub fn Plan::equal(Self, Self) -> Bool

枚举的构造函数

map_plan_kind、reduce_plan_kind、scan_plan_kind 和 map_reduce_plan_kind

返回 Map、Reduce、Scan 和 MapReduce。

pub fn map_plan_kind() -> PlanKind
pub fn reduce_plan_kind() -> PlanKind
pub fn scan_plan_kind() -> PlanKind
pub fn map_reduce_plan_kind() -> PlanKind

i32_type、i64_type、f32_type、f64_type、bytes_type 和 opaque_type

返回 I32、I64、F32、F64、Bytes 和 Opaque(name)。

pub fn i32_type() -> ValueType
pub fn i64_type() -> ValueType
pub fn f32_type() -> ValueType
pub fn f64_type() -> ValueType
pub fn bytes_type() -> ValueType
pub fn opaque_type(String) -> ValueType

sum_reduction、min_reduction、max_reduction 和 custom_reduction

返回 Sum、Min、Max 和 Custom(name)。

pub fn sum_reduction() -> ReductionKernel
pub fn min_reduction() -> ReductionKernel
pub fn max_reduction() -> ReductionKernel
pub fn custom_reduction(String) -> ReductionKernel
test "enumeration constructors" {
  assert_true(@plan.scan_plan_kind() is @plan.Scan)
  debug_inspect(@plan.opaque_type("rgba"), content="Opaque(\"rgba\")")
  debug_inspect(@plan.custom_reduction("xor"), content="Custom(\"xor\")")
}

构建计划

map

构建一个 Map 计划。

pub fn map(String, ValueType, Int, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> Plan

参数依次是标签、元素类型和输入长度;policy 默认为 @shared.native_policy(),ordering 默认为 PreserveInputOrder。归约为 None。

reduce

构建一个带归约内核且顺序为 PreserveInputOrder 的 Reduce 计划。

pub fn reduce(String, ValueType, Int, ReductionKernel, policy? : @shared.ExecutionPolicy) -> Plan

scan

构建一个顺序为 PreserveInputOrder 且没有归约内核的 Scan 计划。

pub fn scan(String, ValueType, Int, policy? : @shared.ExecutionPolicy) -> Plan

原生扫描内核就是前缀和,因此扫描计划不指定内核。

map_reduce

构建一个 MapReduce 计划。

pub fn map_reduce(String, ValueType, Int, ReductionKernel, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> Plan
test "building plans" {
  let policy = @shared.make_execution_policy(worker_count=2, chunk_size=4).unwrap()
  let plan = @plan.reduce("min", @plan.i64_type(), 8, @plan.min_reduction(), policy~)
  assert_true(plan.kind() is @plan.Reduce)
  assert_eq(plan.domain().input_length(), 8)
  assert_eq(plan.reduction(), Some(@plan.min_reduction()))
}

验证

validate

按固定顺序返回计划的所有问题,若没有则返回空数组。

pub fn validate(Plan) -> Array[ValidationIssue]

设 nn 为输入长度,ww 和 cc 为策略中的工作线程数和块大小,检查如下:

问题报告条件
EmptyInputn≤0n \le 0
UnsupportedValueType元素类型不是 I32 或 I64
ChunkSizeDoesNotFitInputc>n>0c > n > 0
WorkerCountExceedsInputw>n>0w > n > 0
UnsupportedBackend策略的后端不是 Native
UnsupportedMode策略的模式不是 Synchronous
MissingReductionKernelReduce 或 MapReduce 计划没有内核
UnsupportedReductionKernelReduce 或 MapReduce 计划的内核是 Custom
UnsupportedOrdering 和 ScanRequiresStableOrderingScan 计划的顺序为 RelaxedOrder;两个问题都会报告

validate 不检查原生内核同样要求的覆盖条件 w≥⌈n/c⌉w \ge \lceil n / c \rceil。

is_runnable

当 validate 没有返回任何问题时返回 true。

pub fn is_runnable(Plan) -> Bool

value_type_is_supported_in_v1、reduction_kernel_is_supported_in_v1 和 ordering_is_valid_for_kind

validate 使用的各条 v1 规则。

pub fn value_type_is_supported_in_v1(ValueType) -> Bool
pub fn reduction_kernel_is_supported_in_v1(ReductionKernel) -> Bool
pub fn ordering_is_valid_for_kind(PlanKind, @shared.OrderingGuarantee) -> Bool

第一个接受 I32 和 I64,第二个接受 Sum、Min 和 Max。第三个接受除 Scan 的 RelaxedOrder 以外的所有顺序。

test "validation" {
  let relaxed = @plan.Plan::new(
    "prefix",
    @plan.scan_plan_kind(),
    @plan.DataDomain::new(@plan.i32_type(), 8),
    ordering=@shared.relaxed_order(),
  )
  let issues = @plan.validate(relaxed)
  assert_eq(issues.length(), 2)
  assert_true(issues[0] is @plan.UnsupportedOrdering(kind=@plan.Scan, ..))
  assert_true(issues[1] is @plan.ScanRequiresStableOrdering)
  assert_true(!@plan.is_runnable(relaxed))
  assert_true(!@plan.reduction_kernel_is_supported_in_v1(@plan.custom_reduction("xor")))
}

谓词

is_map、is_reduce、is_scan 和 is_map_reduce

检查计划的种类。

pub fn is_map(Plan) -> Bool
pub fn is_reduce(Plan) -> Bool
pub fn is_scan(Plan) -> Bool
pub fn is_map_reduce(Plan) -> Bool

has_sum_reduction

当计划的内核为 Some(Sum) 时返回 true。

pub fn has_sum_reduction(Plan) -> Bool

is_chunk_size_too_large、is_scan_requires_stable_ordering、is_unsupported_value_type、is_unsupported_reduction_kernel、is_unsupported_backend、is_unsupported_mode 和 is_unsupported_ordering

判断一个 ValidationIssue 是哪种问题,可配合 Array::any 使用。

pub fn is_chunk_size_too_large(ValidationIssue) -> Bool
pub fn is_scan_requires_stable_ordering(ValidationIssue) -> Bool
pub fn is_unsupported_value_type(ValidationIssue) -> Bool
pub fn is_unsupported_reduction_kernel(ValidationIssue) -> Bool
pub fn is_unsupported_backend(ValidationIssue) -> Bool
pub fn is_unsupported_mode(ValidationIssue) -> Bool
pub fn is_unsupported_ordering(ValidationIssue) -> Bool

is_chunk_size_too_large 匹配 ChunkSizeDoesNotFitInput;其余的匹配同名问题。

test "issue predicates" {
  let policy = @shared.make_execution_policy(chunk_size=64).unwrap()
  let plan = @plan.map("small", @plan.i32_type(), 8, policy~)
  let issues = @plan.validate(plan)
  assert_true(issues.any(@plan.is_chunk_size_too_large))
  assert_true(@plan.is_map(plan) && !@plan.has_sum_reduction(plan))
}

包信息

package_name 和 default_backend_label

返回 "plan" 以及默认后端的标签 "native"。

pub fn package_name() -> String
pub fn default_backend_label() -> String