c
org.apache.spark.sql.execution.python.streaming
PythonForeachWriter
Companion object PythonForeachWriter
class PythonForeachWriter extends ForeachWriter[UnsafeRow]
The class proceeds as follows:
- Rows streamed through a
process()call on the org.apache.spark.sql.execution.streaming.QueryExecutionThread are buffered in theUnsafeRowBuffer. - The WriterThread streams the buffered data to the Python worker. - Once the streaming query ends, close() is called which signals the buffer to mark the end of streaming input. The streaming query execution thread waits for the WriterThread to complete and throws any exceptions seen by the WriterThread.
Linear Supertypes
Ordering
- Alphabetic
- By Inheritance
Inherited
- PythonForeachWriter
- ForeachWriter
- Serializable
- AnyRef
- Any
- Hide All
- Show All
Visibility
- Public
- Protected
Instance Constructors
- new PythonForeachWriter(func: PythonFunction, schema: StructType)
Value Members
- final def !=(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
- final def ##: Int
- Definition Classes
- AnyRef → Any
- final def ==(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
- final def asInstanceOf[T0]: T0
- Definition Classes
- Any
- def clone(): AnyRef
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
- def close(errorOrNull: Throwable): Unit
Waits for the writer thread to finish evaluating the Python function.
Waits for the writer thread to finish evaluating the Python function. Throws any exceptions seen by the writer thread.
- Definition Classes
- PythonForeachWriter → ForeachWriter
- final def eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
- def equals(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef → Any
- final def getClass(): Class[_ <: AnyRef]
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- def hashCode(): Int
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- final def isInstanceOf[T0]: Boolean
- Definition Classes
- Any
- final def ne(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
- final def notify(): Unit
- Definition Classes
- AnyRef
- Annotations
- @IntrinsicCandidate() @native()
- final def notifyAll(): Unit
- Definition Classes
- AnyRef
- Annotations
- @IntrinsicCandidate() @native()
- def open(partitionId: Long, version: Long): Boolean
- Definition Classes
- PythonForeachWriter → ForeachWriter
- def process(value: UnsafeRow): Unit
- Definition Classes
- PythonForeachWriter → ForeachWriter
- final def synchronized[T0](arg0: => T0): T0
- Definition Classes
- AnyRef
- def toString(): String
- Definition Classes
- AnyRef → Any
- final def wait(arg0: Long, arg1: Int): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException])
- final def wait(arg0: Long): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException]) @native()
- final def wait(): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException])
Deprecated Value Members
- def finalize(): Unit
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.Throwable]) @Deprecated
- Deprecated
(Since version 9)