Packages

c

org.apache.spark.sql.execution.python.streaming

TransformWithStateInPySparkPythonPreInitRunner

class TransformWithStateInPySparkPythonPreInitRunner extends StreamingPythonRunner with TransformWithStateInPySparkPythonRunnerUtils with Logging

TransformWithStateInPySpark driver side Python runner. Similar as executor side runner, will start a new daemon thread on the Python runner to run state server.

Linear Supertypes
TransformWithStateInPySparkPythonRunnerUtils, StreamingPythonRunner, Logging, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. TransformWithStateInPySparkPythonPreInitRunner
  2. TransformWithStateInPySparkPythonRunnerUtils
  3. StreamingPythonRunner
  4. Logging
  5. AnyRef
  6. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. Protected

Instance Constructors

  1. new TransformWithStateInPySparkPythonPreInitRunner(func: PythonFunction, workerModule: String, groupingKeySchema: StructType, processorHandleImpl: DriverStatefulProcessorHandleImpl)

Type Members

  1. implicit class LogStringContext extends AnyRef
    Definition Classes
    Logging
  2. class StreamingPythonRunnerInitializationCommunicationException extends SparkPythonException
    Definition Classes
    StreamingPythonRunner
  3. class StreamingPythonRunnerInitializationException extends SparkPythonException
    Definition Classes
    StreamingPythonRunner
  4. class StreamingPythonRunnerInitializationTimeoutException extends SparkPythonException
    Definition Classes
    StreamingPythonRunner

Value Members

  1. final def !=(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  2. final def ##: Int
    Definition Classes
    AnyRef → Any
  3. final def ==(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  4. def MDC(key: LogKey, value: Any): MDC
    Attributes
    protected
    Definition Classes
    Logging
  5. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  6. val authSocketTimeout: Long
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  7. val bufferSize: Int
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  8. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
  9. def closeServerSocketChannelSilently(stateServerSocket: ServerSocketChannel): Unit
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  10. val envVars: Map[String, String]
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  11. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  12. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  13. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  14. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  15. def init(): (DataOutputStream, DataInputStream)
    Definition Classes
    TransformWithStateInPySparkPythonPreInitRunner → StreamingPythonRunner
  16. def initStateServer(): Unit
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  17. def initializeLogIfNecessary(isInterpreter: Boolean, silent: Boolean): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  18. def initializeLogIfNecessary(isInterpreter: Boolean): Unit
    Attributes
    protected
    Definition Classes
    Logging
  19. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  20. def isTraceEnabled(): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  21. val isUnixDomainSock: Boolean
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  22. def isWorkerStopped(): Option[Boolean]
    Definition Classes
    StreamingPythonRunner
  23. def log: Logger
    Attributes
    protected
    Definition Classes
    Logging
  24. def logBasedOnLevel(level: Level)(f: => MessageWithContext): Unit
    Attributes
    protected
    Definition Classes
    Logging
  25. def logDebug(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  26. def logDebug(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  27. def logDebug(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  28. def logDebug(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  29. def logError(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  30. def logError(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  31. def logError(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  32. def logError(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  33. def logInfo(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  34. def logInfo(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  35. def logInfo(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  36. def logInfo(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  37. def logName: String
    Attributes
    protected
    Definition Classes
    Logging
  38. def logTrace(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  39. def logTrace(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  40. def logTrace(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  41. def logTrace(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  42. def logWarning(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  43. def logWarning(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  44. def logWarning(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  45. def logWarning(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  46. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  47. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  48. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  49. def process(): Unit
  50. val pythonExec: String
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  51. val pythonVer: String
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  52. var pythonWorker: Option[PythonWorker]
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  53. var pythonWorkerFactory: Option[PythonWorkerFactory]
    Attributes
    protected
    Definition Classes
    StreamingPythonRunner
  54. val sqlConf: SQLConf
    Attributes
    protected
  55. val stateServerSocket: ServerSocketChannel
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  56. val stateServerSocketPath: String
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  57. val stateServerSocketPort: Int
    Attributes
    protected
    Definition Classes
    TransformWithStateInPySparkPythonRunnerUtils
  58. def stop(): Unit
    Definition Classes
    TransformWithStateInPySparkPythonPreInitRunner → StreamingPythonRunner
  59. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  60. def toString(): String
    Definition Classes
    AnyRef → Any
  61. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  62. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException]) @native()
  63. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  64. def withLogContext(context: Map[String, String])(body: => Unit): Unit
    Attributes
    protected
    Definition Classes
    Logging

Deprecated Value Members

  1. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.Throwable]) @Deprecated
    Deprecated

    (Since version 9)

Inherited from StreamingPythonRunner

Inherited from Logging

Inherited from AnyRef

Inherited from Any

Ungrouped