beam: github.com/apache/beam/sdks/go/pkg/beam/core/runtime/pipelinex Index | Files

package pipelinex

import "github.com/apache/beam/sdks/go/pkg/beam/core/runtime/pipelinex"

Package pipelinex contains utilities for manipulating Beam proto pipelines. The utilities generally uses shallow copies and do not mutate their inputs.

Index

Package Files

clone.go replace.go util.go

func Bounded Uses

func Bounded(p *pb.Pipeline) bool

Bounded returns true iff all PCollections are bounded.

func ContainerImages Uses

func ContainerImages(p *pb.Pipeline) []string

ContainerImages returns the set of container images used in the given pipeline.

func Normalize Uses

func Normalize(p *pb.Pipeline) (*pb.Pipeline, error)

Normalize recomputes derivative information in the pipeline, such as roots and input/output for composite transforms. It also ensures that unique names are so and topologically sorts each subtransform list.

func ShallowCloneFunctionSpec Uses

func ShallowCloneFunctionSpec(p *pb.FunctionSpec) *pb.FunctionSpec

ShallowCloneFunctionSpec makes a shallow copy of the given FunctionSpec.

func ShallowClonePTransform Uses

func ShallowClonePTransform(t *pb.PTransform) *pb.PTransform

ShallowClonePTransform makes a shallow copy of the given PTransform.

func ShallowCloneParDoPayload Uses

func ShallowCloneParDoPayload(p *pb.ParDoPayload) *pb.ParDoPayload

ShallowCloneParDoPayload makes a shallow copy of the given ParDoPayload.

func ShallowCloneSideInput Uses

func ShallowCloneSideInput(p *pb.SideInput) *pb.SideInput

ShallowCloneSideInput makes a shallow copy of the given SideInput.

func TopologicalSort Uses

func TopologicalSort(xforms map[string]*pb.PTransform, ids []string) []string

TopologicalSort returns a topologically sorted list of the given ids, generally from the same scope/composite. Assumes acyclic graph.

func TrimCoders Uses

func TrimCoders(coders map[string]*pb.Coder, ids ...string) map[string]*pb.Coder

TrimCoders returns the transitive closure of the given coders ids.

func Update Uses

func Update(p *pb.Pipeline, values *pb.Components) (*pb.Pipeline, error)

Update merges a pipeline with the given components, which may add, replace or delete its values. It returns the merged pipeline. The input is not modified.

Package pipelinex imports 6 packages (graph) and is imported by 2 packages. Updated 2020-01-18. Refresh now. Tools for package owners.