Class WriteFiles<UserT,DestinationT,OutputT>

java.lang.Object
org.apache.beam.sdk.transforms.PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
org.apache.beam.sdk.io.WriteFiles<UserT,DestinationT,OutputT>
All Implemented Interfaces:
Serializable, HasDisplayData

public abstract class WriteFiles<UserT,DestinationT,OutputT> extends PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
A PTransform that writes to a FileBasedSink. A write begins with a sequential global initialization of a sink, followed by a parallel write, and ends with a sequential finalization of the write. The output of a write is PDone.

By default, every bundle in the input PCollection will be processed by a FileBasedSink.WriteOperation, so the number of output will vary based on runner behavior, though at least 1 output will always be produced. The exact parallelism of the write stage can be controlled using withNumShards(int), typically used to control how many files are produced or to globally limit the number of workers connecting to an external service. However, this option can often hurt performance: it adds an additional GroupByKey to the pipeline.

Example usage with runner-determined sharding:

p.apply(WriteFiles.to(new MySink(...)));

Example usage with a fixed number of shards:

p.apply(WriteFiles.to(new MySink(...)).withNumShards(3));
See Also:
  • Field Details

    • CONCRETE_CLASS

      @Internal public static final Class<? extends WriteFiles> CONCRETE_CLASS
      For internal use by runners.
    • FILE_TRIGGERING_RECORD_COUNT

      public static final int FILE_TRIGGERING_RECORD_COUNT
      See Also:
    • FILE_TRIGGERING_BYTE_COUNT

      public static final int FILE_TRIGGERING_BYTE_COUNT
      See Also:
    • FILE_TRIGGERING_RECORD_BUFFERING_DURATION

      public static final Duration FILE_TRIGGERING_RECORD_BUFFERING_DURATION
  • Constructor Details

    • WriteFiles

      public WriteFiles()
  • Method Details

    • to

      public static <UserT, DestinationT, OutputT> WriteFiles<UserT,DestinationT,OutputT> to(FileBasedSink<UserT,DestinationT,OutputT> sink)
      Creates a WriteFiles transform that writes to the given FileBasedSink, letting the runner control how many different shards are produced.
    • getSink

      public abstract FileBasedSink<UserT,DestinationT,OutputT> getSink()
    • getComputeNumShards

      public abstract @Nullable PTransform<PCollection<UserT>,PCollectionView<Integer>> getComputeNumShards()
    • getNumShardsProvider

      public abstract @Nullable ValueProvider<Integer> getNumShardsProvider()
    • getWindowedWrites

      public abstract boolean getWindowedWrites()
    • getWithAutoSharding

      public abstract boolean getWithAutoSharding()
    • getShardingFunction

      public abstract @Nullable ShardingFunction<UserT,DestinationT> getShardingFunction()
    • getBadRecordErrorHandler

      public abstract ErrorHandler<BadRecord,?> getBadRecordErrorHandler()
    • getBadRecordRouter

      public abstract BadRecordRouter getBadRecordRouter()
    • getAdditionalInputs

      public Map<TupleTag<?>,PValue> getAdditionalInputs()
      Description copied from class: PTransform
      Returns all PValues that are consumed as inputs to this PTransform that are independent of the expansion of the PTransform within PTransform.expand(PInput).

      For example, this can contain any side input consumed by this PTransform.

      Overrides:
      getAdditionalInputs in class PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
    • withNumShards

      public WriteFiles<UserT,DestinationT,OutputT> withNumShards(int numShards)
      Returns a new WriteFiles that will write to the current FileBasedSink using the specified number of shards.

      This option should be used sparingly as it can hurt performance. See WriteFiles for more information.

      A value less than or equal to 0 will be equivalent to the default behavior of runner-determined sharding.

    • withNumShards

      public WriteFiles<UserT,DestinationT,OutputT> withNumShards(ValueProvider<Integer> numShardsProvider)
      Returns a new WriteFiles that will write to the current FileBasedSink using the ValueProvider specified number of shards.

      This option should be used sparingly as it can hurt performance. See WriteFiles for more information.

    • withMaxNumWritersPerBundle

      public WriteFiles<UserT,DestinationT,OutputT> withMaxNumWritersPerBundle(int maxNumWritersPerBundle)
      Set the maximum number of writers kept open in a bundle before spilling to shuffle (or evicting the least recently used open writer if withEvictWritersWhenFull() is enabled).

      Trade-offs: A higher value here can cause more worker memory consumption (since each open writer maintains an in-memory write buffer), but reduces the cost of shuffling spilled records (or reduces how frequently writers are closed and evicted when withEvictWritersWhenFull() is enabled). A lower value reduces peak memory consumption per bundle at the cost of either more records spilled to shuffle or more frequent writer evictions (resulting in smaller output files).

      Writer Limit & Overflow Trade-off Matrix (for withRunnerDeterminedSharding()):

      Configuration Behavior when maxNumWritersPerBundle is reached Worker Memory Consumption Shuffle Cost Output File Size / Count
      Default (Spill to Shuffle)
      maxNumWritersPerBundle > 0,
      evictWritersWhenFull = false
      Keeps first N writers open; spills remaining records to a GroupByKey shuffle stage Bounded (<= N buffers per bundle) High if many records spill across shuffle Fewer, larger files
      LRU Writer Eviction
      maxNumWritersPerBundle > 0,
      evictWritersWhenFull = true
      Flushes and closes the least recently used (LRU) open writer to open a new writer inline Bounded (<= N buffers per bundle) None (no shuffle stage for unwritten records) May produce more/smaller files (minimal if input is ordered by destination, high if random)
      No Spilling
      withNoSpilling() (maxNumWritersPerBundle = -1)
      Opens a new writer for every destination in the bundle without limit Unbounded (risk of OOM with many destinations) None (no shuffle stage for unwritten records) Fewer, larger files (1 file per destination per bundle)

      Note that value provided here cannot exceed the default value (DEFAULT_MAX_NUM_WRITERS_PER_BUNDLE).

    • withEvictWritersWhenFull

      public WriteFiles<UserT,DestinationT,OutputT> withEvictWritersWhenFull()
      Returns a new WriteFiles that evicts the least recently used open writer in the bundle (LRU order, by flushing and closing it) instead of spilling unwritten records to shuffle when getMaxNumWritersPerBundle() is reached.

      Trade-offs: Setting this to true avoids the cost of shuffling records while keeping concurrent writer memory consumption bounded by getMaxNumWritersPerBundle(), but may lead to smaller and more numerous output files since evicted writers are closed before the end of the bundle. See withMaxNumWritersPerBundle(int) for the full trade-off matrix.

      Warning: This option should only be used when the input PCollection elements within a bundle are already grouped or ordered by writer keys (destination/window/pane), such that consecutive records belong to the same destination. If the input PCollection rows arrive in random order across more destinations than getMaxNumWritersPerBundle(), writers will be repeatedly closed and reopened, creating too many small files.

      This option only applies to writes withRunnerDeterminedSharding().

    • withEvictWritersWhenFull

      public WriteFiles<UserT,DestinationT,OutputT> withEvictWritersWhenFull(boolean evictWritersWhenFull)
      Set this sink to evict the least recently used open writer in the bundle (LRU order, by flushing and closing it) when getMaxNumWritersPerBundle() is reached, instead of spilling unwritten records to shuffle.

      Trade-offs: Setting this to true avoids the cost of shuffling records while keeping concurrent writer memory consumption bounded by getMaxNumWritersPerBundle(), but may lead to smaller and more numerous output files since evicted writers are closed before the end of the bundle. Setting this to false (default) preserves larger output files by spilling excess records to a shuffle stage. See withMaxNumWritersPerBundle(int) for the full trade-off matrix.

      Warning: This option should only be used when the input PCollection elements within a bundle are already grouped or ordered by writer keys (destination/window/pane), such that consecutive records belong to the same destination. If the input PCollection rows arrive in random order across more destinations than getMaxNumWritersPerBundle(), writers will be repeatedly closed and reopened, creating too many small files.

      This option only applies to writes withRunnerDeterminedSharding().

    • withSkipIfEmpty

      public WriteFiles<UserT,DestinationT,OutputT> withSkipIfEmpty(boolean skipIfEmpty)
      Set this sink to skip writing any files if the PCollection is empty.
    • withBatchSize

      public WriteFiles<UserT,DestinationT,OutputT> withBatchSize(@Nullable Integer batchSize)
      Returns a new WriteFiles that will batch the input records using specified batch size. The default value is FILE_TRIGGERING_RECORD_COUNT.

      This option is used only for writing unbounded data with auto-sharding.

    • withBatchSizeBytes

      public WriteFiles<UserT,DestinationT,OutputT> withBatchSizeBytes(@Nullable Integer batchSizeBytes)
      Returns a new WriteFiles that will batch the input records using specified batch size in bytes. The default value is FILE_TRIGGERING_BYTE_COUNT.

      This option is used only for writing unbounded data with auto-sharding.

    • withBatchMaxBufferingDuration

      public WriteFiles<UserT,DestinationT,OutputT> withBatchMaxBufferingDuration(@Nullable Duration batchMaxBufferingDuration)
      Returns a new WriteFiles that will batch the input records using specified max buffering duration. The default value is FILE_TRIGGERING_RECORD_BUFFERING_DURATION.

      This option is used only for writing unbounded data with auto-sharding.

    • withSideInputs

      public WriteFiles<UserT,DestinationT,OutputT> withSideInputs(List<PCollectionView<?>> sideInputs)
    • withSharding

      Returns a new WriteFiles that will write to the current FileBasedSink using the specified PTransform to compute the number of shards.

      This option should be used sparingly as it can hurt performance. See WriteFiles for more information.

    • withRunnerDeterminedSharding

      public WriteFiles<UserT,DestinationT,OutputT> withRunnerDeterminedSharding()
      Returns a new WriteFiles that will write to the current FileBasedSink with runner-determined sharding.
    • withAutoSharding

      public WriteFiles<UserT,DestinationT,OutputT> withAutoSharding()
    • withShardingFunction

      public WriteFiles<UserT,DestinationT,OutputT> withShardingFunction(ShardingFunction<UserT,DestinationT> shardingFunction)
      Returns a new WriteFiles that will write to the current FileBasedSink using the specified sharding function to assign shard for inputs.
    • withWindowedWrites

      public WriteFiles<UserT,DestinationT,OutputT> withWindowedWrites()
      Returns a new WriteFiles that writes preserves windowing on it's input.

      If this option is not specified, windowing and triggering are replaced by GlobalWindows and DefaultTrigger.

      If there is no data for a window, no output shards will be generated for that window. If a window triggers multiple times, then more than a single output shard might be generated multiple times; it's up to the sink implementation to keep these output shards unique.

      This option can only be used if withNumShards(int) is also set to a positive value.

    • withNoSpilling

      public WriteFiles<UserT,DestinationT,OutputT> withNoSpilling()
      Returns a new WriteFiles that writes all data without spilling, simplifying the pipeline. This option should not be used with withMaxNumWritersPerBundle(int) and it will eliminate this limit possibly causing many writers to be opened. Use with caution.

      This option only applies to writes withRunnerDeterminedSharding().

    • withSkipIfEmpty

      public WriteFiles<UserT,DestinationT,OutputT> withSkipIfEmpty()
    • withBadRecordErrorHandler

      public WriteFiles<UserT,DestinationT,OutputT> withBadRecordErrorHandler(ErrorHandler<BadRecord,?> errorHandler)
    • validate

      public void validate(PipelineOptions options)
      Description copied from class: PTransform
      Called before running the Pipeline to verify this transform is fully and correctly specified.

      By default, does nothing.

      Overrides:
      validate in class PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
    • expand

      Description copied from class: PTransform
      Override this method to specify how this PTransform should be expanded on the given InputT.

      NOTE: This method should not be called directly. Instead apply the PTransform should be applied to the InputT using the apply method.

      Composite transforms, which are defined in terms of other transforms, should return the output of one of the composed transforms. Non-composite transforms, which do not apply any transforms internally, should return a new unbound output and register evaluators (via backend-specific registration methods).

      Specified by:
      expand in class PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
    • populateDisplayData

      public void populateDisplayData(DisplayData.Builder builder)
      Description copied from class: PTransform
      Register display data for the given transform or component.

      populateDisplayData(DisplayData.Builder) is invoked by Pipeline runners to collect display data via DisplayData.from(HasDisplayData). Implementations may call super.populateDisplayData(builder) in order to register display data in the current namespace, but should otherwise use subcomponent.populateDisplayData(builder) to use the namespace of the subcomponent.

      By default, does not register any display data. Implementors may override this method to provide their own display data.

      Specified by:
      populateDisplayData in interface HasDisplayData
      Overrides:
      populateDisplayData in class PTransform<PCollection<UserT>,WriteFilesResult<DestinationT>>
      Parameters:
      builder - The builder to populate with display data.
      See Also: