Modifier and Type | Class and Description |
---|---|
class |
ProcessPythonEnvironmentManager
The ProcessPythonEnvironmentManager is used to prepare the working dir of python UDF worker and
create ProcessPythonEnvironment object of Beam Fn API.
|
Modifier and Type | Method and Description |
---|---|
protected PythonEnvironmentManager |
AbstractPythonFunctionOperator.createPythonEnvironmentManager() |
Constructor and Description |
---|
BeamDataStreamPythonFunctionRunner(String taskName,
PythonEnvironmentManager environmentManager,
TypeInformation inputType,
TypeInformation outputType,
String functionUrn,
FlinkFnApi.UserDefinedDataStreamFunction userDefinedDataStreamFunction,
String coderUrn,
Map<String,String> jobOptions,
FlinkMetricContainer flinkMetricContainer,
KeyedStateBackend stateBackend,
TypeSerializer keySerializer,
MemoryManager memoryManager,
double managedMemoryFraction) |
BeamPythonFunctionRunner(String taskName,
PythonEnvironmentManager environmentManager,
String functionUrn,
Map<String,String> jobOptions,
FlinkMetricContainer flinkMetricContainer,
KeyedStateBackend keyedStateBackend,
TypeSerializer keySerializer,
MemoryManager memoryManager,
double managedMemoryFraction) |
Modifier and Type | Method and Description |
---|---|
protected PythonEnvironmentManager |
AbstractPythonStatelessFunctionFlatMap.createPythonEnvironmentManager() |
Constructor and Description |
---|
BeamTableStatefulPythonFunctionRunner(String taskName,
PythonEnvironmentManager environmentManager,
RowType inputType,
RowType outputType,
String functionUrn,
FlinkFnApi.UserDefinedAggregateFunctions userDefinedFunctions,
String coderUrn,
Map<String,String> jobOptions,
FlinkMetricContainer flinkMetricContainer,
KeyedStateBackend keyedStateBackend,
TypeSerializer keySerializer,
MemoryManager memoryManager,
double managedMemoryFraction) |
BeamTableStatelessPythonFunctionRunner(String taskName,
PythonEnvironmentManager environmentManager,
RowType inputType,
RowType outputType,
String functionUrn,
FlinkFnApi.UserDefinedFunctions userDefinedFunctions,
String coderUrn,
Map<String,String> jobOptions,
FlinkMetricContainer flinkMetricContainer,
MemoryManager memoryManager,
double managedMemoryFraction) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.