Class PythonTransforms

java.lang.Object
com.hazelcast.jet.python.PythonTransforms

public final class PythonTransforms extends Object
Transforms which allow the user to call Python user-defined functions from inside a Jet pipeline.

See also PythonExtension for a fluent and modernized API.

Since:
Jet 4.0
  • Method Summary

    Modifier and Type
    Method
    Description
    static <K> com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>>
    mapUsingPython(com.hazelcast.function.FunctionEx<? super String,? extends K> keyFn, PythonServiceConfig cfg)
    Deprecated, for removal: This API element is subject to removal in a future version.
    Jet now has first-class support for data rebalancing, see GeneralStage.rebalance() and GeneralStage.rebalance(FunctionEx).
    static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>>
    A stage-transforming method that adds a "map using Python" pipeline stage.
    static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>>
    mapUsingPython(PythonServiceConfig cfg, int maxBatchSize)
    A stage-transforming method that adds a "map using Python" pipeline stage.
    static <K> com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>>
    mapUsingPythonBatch(com.hazelcast.function.FunctionEx<? super String,? extends K> keyFn, PythonServiceConfig cfg)
    Deprecated, for removal: This API element is subject to removal in a future version.
    Jet now has first-class support for data rebalancing, see GeneralStage.rebalance() and GeneralStage.rebalance(FunctionEx).
    static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>>
    A stage-transforming method that adds a "map using Python" pipeline stage.
    static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>>
    mapUsingPythonBatch(PythonServiceConfig cfg, int maxBatchSize)
    A stage-transforming method that adds a "map using Python" pipeline stage.

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Method Details

    • mapUsingPython

      @Nonnull public static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>> mapUsingPython(@Nonnull PythonServiceConfig cfg)
      A stage-transforming method that adds a "map using Python" pipeline stage. Use it with stage.apply(PythonService.mapUsingPython(pyConfig)). See PythonServiceConfig for more details.
    • mapUsingPython

      @Nonnull public static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>> mapUsingPython(@Nonnull PythonServiceConfig cfg, int maxBatchSize)
      A stage-transforming method that adds a "map using Python" pipeline stage. Use it with stage.apply(PythonService.mapUsingPython(pyConfig)). See PythonServiceConfig for more details.

      The maxBatchSize may be used to limit the size of a single request to the Python service.

      Parameters:
      cfg - configuration for the Python service
      maxBatchSize - the maximum size of a batch for a single request
    • mapUsingPython

      @Deprecated(since="5.0", forRemoval=true) @Nonnull public static <K> com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.StreamStage<String>,com.hazelcast.jet.pipeline.StreamStage<String>> mapUsingPython(@Nonnull com.hazelcast.function.FunctionEx<? super String,? extends K> keyFn, @Nonnull PythonServiceConfig cfg)
      Deprecated, for removal: This API element is subject to removal in a future version.
      Jet now has first-class support for data rebalancing, see GeneralStage.rebalance() and GeneralStage.rebalance(FunctionEx).
      A stage-transforming method that adds a partitioned "map using Python" pipeline stage. It applies partitioning using the supplied keyFn. You need partitioning if your input stream comes from a non-distributed data source (all data coming in on a single cluster member), in order to distribute the Python work across the whole cluster.

      Use it like this: stage.apply(PythonService.mapUsingPython(keyFn, pyConfig)). See PythonServiceConfig for more details.

    • mapUsingPythonBatch

      @Nonnull public static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>> mapUsingPythonBatch(@Nonnull PythonServiceConfig cfg)
      A stage-transforming method that adds a "map using Python" pipeline stage. Use it with stage.apply(PythonService.mapUsingPythonBatch(pyConfig)). See PythonServiceConfig for more details.
    • mapUsingPythonBatch

      @Nonnull public static com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>> mapUsingPythonBatch(@Nonnull PythonServiceConfig cfg, int maxBatchSize)
      A stage-transforming method that adds a "map using Python" pipeline stage. Use it with stage.apply(PythonService.mapUsingPythonBatch(pyConfig)). See PythonServiceConfig for more details.

      The maxBatchSize may be used to limit the size of a single request to the Python service.

      Parameters:
      cfg - configuration for the Python service
      maxBatchSize - the maximum size of a batch for a single request
    • mapUsingPythonBatch

      @Nonnull @Deprecated(since="5.0", forRemoval=true) public static <K> com.hazelcast.function.FunctionEx<com.hazelcast.jet.pipeline.BatchStage<String>,com.hazelcast.jet.pipeline.BatchStage<String>> mapUsingPythonBatch(@Nonnull com.hazelcast.function.FunctionEx<? super String,? extends K> keyFn, @Nonnull PythonServiceConfig cfg)
      Deprecated, for removal: This API element is subject to removal in a future version.
      Jet now has first-class support for data rebalancing, see GeneralStage.rebalance() and GeneralStage.rebalance(FunctionEx).
      A stage-transforming method that adds a partitioned "map using Python" pipeline stage. It applies partitioning using the supplied keyFn. You need partitioning if your input stream comes from a non-distributed data source (all data coming in on a single cluster member), in order to distribute the Python work across the whole cluster.

      Use it like this: stage.apply(PythonService.mapUsingPythonBatch(keyFn, pyConfig)). See PythonServiceConfig for more details.