schrodinger.seam.testing.stu_tests.jobfs_compression module

Test that a worker can write a gzip-compressed file under jobfs:// and the driver correctly reads it back after _processJobfsFiles xcopies the file into the driver’s jobfs directory.

This guards against the silent-corruption regression where a .gz-named jobfs file held plaintext because JobFileSystem ignored compression_type.

Usage:

$SCHRODINGER/run python3 -m schrodinger.seam.testing.stu_tests.jobfs_compression
class schrodinger.seam.testing.stu_tests.jobfs_compression.WriteCompressedJobFS(*unused_args, **unused_kwargs)

Bases: DoFn

Write a gzip-compressed jobfs file for each element.

process(element)

Method to use for processing elements.

This is invoked by DoFnRunner for each element of a input PCollection.

The following parameters can be used as default values on process arguments to indicate that a DoFn accepts the corresponding parameters. For example, a DoFn might accept the element and its timestamp with the following signature:

def process(element=DoFn.ElementParam, timestamp=DoFn.TimestampParam):
  ...

The full set of parameters is:

  • DoFn.ElementParam: element to be processed, should not be mutated.

  • DoFn.SideInputParam: a side input that may be used when processing.

  • DoFn.TimestampParam: timestamp of the input element.

  • DoFn.WindowParam: Window the input element belongs to.

  • DoFn.TimerParam: a userstate.RuntimeTimer object defined by the spec of the parameter.

  • DoFn.StateParam: a userstate.RuntimeState object defined by the spec of the parameter.

  • DoFn.KeyParam: key associated with the element.

  • DoFn.RestrictionParam: an iobase.RestrictionTracker will be provided here to allow treatment as a Splittable DoFn. The restriction tracker will be derived from the restriction provider in the parameter.

  • DoFn.WatermarkEstimatorParam: a function that can be used to track output watermark of Splittable DoFn implementations.

  • DoFn.BundleContextParam: allows a shared context manager to be used per bundle

  • DoFn.SetupContextParam: allows a shared context manager to be used per DoFn

Args:

element: The element to be processed *args: side inputs **kwargs: other keyword arguments.

Returns:

An Iterable of output elements or None.

class schrodinger.seam.testing.stu_tests.jobfs_compression.ReadCompressedJobFS(*unused_args, **unused_kwargs)

Bases: DoFn

Read the per-element jobfs file and yield element:plaintext.

process(element)

Method to use for processing elements.

This is invoked by DoFnRunner for each element of a input PCollection.

The following parameters can be used as default values on process arguments to indicate that a DoFn accepts the corresponding parameters. For example, a DoFn might accept the element and its timestamp with the following signature:

def process(element=DoFn.ElementParam, timestamp=DoFn.TimestampParam):
  ...

The full set of parameters is:

  • DoFn.ElementParam: element to be processed, should not be mutated.

  • DoFn.SideInputParam: a side input that may be used when processing.

  • DoFn.TimestampParam: timestamp of the input element.

  • DoFn.WindowParam: Window the input element belongs to.

  • DoFn.TimerParam: a userstate.RuntimeTimer object defined by the spec of the parameter.

  • DoFn.StateParam: a userstate.RuntimeState object defined by the spec of the parameter.

  • DoFn.KeyParam: key associated with the element.

  • DoFn.RestrictionParam: an iobase.RestrictionTracker will be provided here to allow treatment as a Splittable DoFn. The restriction tracker will be derived from the restriction provider in the parameter.

  • DoFn.WatermarkEstimatorParam: a function that can be used to track output watermark of Splittable DoFn implementations.

  • DoFn.BundleContextParam: allows a shared context manager to be used per bundle

  • DoFn.SetupContextParam: allows a shared context manager to be used per DoFn

Args:

element: The element to be processed *args: side inputs **kwargs: other keyword arguments.

Returns:

An Iterable of output elements or None.

schrodinger.seam.testing.stu_tests.jobfs_compression.main()