Packages

c

pl.touk.nussknacker.engine.kafka.source.flink

FlinkKafkaConsumerHandlingExceptions

class FlinkKafkaConsumerHandlingExceptions[T] extends FlinkKafkaConsumer[T] with LazyLogging

Annotations
@silent("deprecated")
Linear Supertypes
LazyLogging, FlinkKafkaConsumer[T], FlinkKafkaConsumerBase[T], CheckpointedFunction, ResultTypeQueryable[T], CheckpointListener, RichParallelSourceFunction[T], ParallelSourceFunction[T], SourceFunction[T], AbstractRichFunction, RichFunction, Function, Serializable, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. FlinkKafkaConsumerHandlingExceptions
  2. LazyLogging
  3. FlinkKafkaConsumer
  4. FlinkKafkaConsumerBase
  5. CheckpointedFunction
  6. ResultTypeQueryable
  7. CheckpointListener
  8. RichParallelSourceFunction
  9. ParallelSourceFunction
  10. SourceFunction
  11. AbstractRichFunction
  12. RichFunction
  13. Function
  14. Serializable
  15. AnyRef
  16. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. Protected

Instance Constructors

  1. new FlinkKafkaConsumerHandlingExceptions(topics: List[String], deserializationSchema: FlinkDeserializationSchemaWrapper[T], props: Properties, exceptionHandlerPreparer: (RuntimeContext) => ExceptionHandler, convertToEngineRuntimeContext: (RuntimeContext) => EngineRuntimeContext, nodeId: NodeId)

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. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  5. def assignTimestampsAndWatermarks(arg0: WatermarkStrategy[T]): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  6. def cancel(): Unit
    Definition Classes
    FlinkKafkaConsumerBase → SourceFunction
  7. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @HotSpotIntrinsicCandidate() @native()
  8. def close(): Unit
    Definition Classes
    FlinkKafkaConsumerHandlingExceptions → FlinkKafkaConsumerBase → AbstractRichFunction → RichFunction
  9. def createFetcher(arg0: SourceContext[T], arg1: Map[KafkaTopicPartition, Long], arg2: SerializedValue[WatermarkStrategy[T]], arg3: StreamingRuntimeContext, arg4: OffsetCommitMode, arg5: MetricGroup, arg6: Boolean): AbstractFetcher[T, _ <: AnyRef]
    Attributes
    protected[org.apache.flink.streaming.connectors.kafka]
    Definition Classes
    FlinkKafkaConsumer → FlinkKafkaConsumerBase
    Annotations
    @throws(classOf[java.lang.Exception])
  10. def createPartitionDiscoverer(arg0: KafkaTopicsDescriptor, arg1: Int, arg2: Int): AbstractPartitionDiscoverer
    Attributes
    protected[org.apache.flink.streaming.connectors.kafka]
    Definition Classes
    FlinkKafkaConsumer → FlinkKafkaConsumerBase
  11. def disableFilterRestoredPartitionsWithSubscribedTopics(): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  12. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  13. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  14. var exceptionHandler: ExceptionHandler
    Attributes
    protected
  15. var exceptionPurposeContextIdGenerator: ContextIdGenerator
    Attributes
    protected
  16. def fetchOffsetsWithTimestamp(arg0: Collection[KafkaTopicPartition], arg1: Long): Map[KafkaTopicPartition, Long]
    Attributes
    protected[org.apache.flink.streaming.connectors.kafka]
    Definition Classes
    FlinkKafkaConsumer → FlinkKafkaConsumerBase
  17. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @HotSpotIntrinsicCandidate() @native()
  18. def getEnableCommitOnCheckpoints(): Boolean
    Definition Classes
    FlinkKafkaConsumerBase
  19. def getIsAutoCommitEnabled(): Boolean
    Attributes
    protected[org.apache.flink.streaming.connectors.kafka]
    Definition Classes
    FlinkKafkaConsumer → FlinkKafkaConsumerBase
  20. def getIterationRuntimeContext(): IterationRuntimeContext
    Definition Classes
    AbstractRichFunction → RichFunction
  21. def getProducedType(): TypeInformation[T]
    Definition Classes
    FlinkKafkaConsumerBase → ResultTypeQueryable
  22. def getRuntimeContext(): RuntimeContext
    Definition Classes
    AbstractRichFunction → RichFunction
  23. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @HotSpotIntrinsicCandidate() @native()
  24. final def initializeState(arg0: FunctionInitializationContext): Unit
    Definition Classes
    FlinkKafkaConsumerBase → CheckpointedFunction
    Annotations
    @throws(classOf[java.lang.Exception])
  25. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  26. lazy val logger: Logger
    Attributes
    protected
    Definition Classes
    LazyLogging
    Annotations
    @transient()
  27. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  28. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @HotSpotIntrinsicCandidate() @native()
  29. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @HotSpotIntrinsicCandidate() @native()
  30. def notifyCheckpointAborted(arg0: Long): Unit
    Definition Classes
    FlinkKafkaConsumerBase → CheckpointListener
  31. final def notifyCheckpointComplete(arg0: Long): Unit
    Definition Classes
    FlinkKafkaConsumerBase → CheckpointListener
    Annotations
    @throws(classOf[java.lang.Exception])
  32. def open(parameters: Configuration): Unit
    Definition Classes
    FlinkKafkaConsumerHandlingExceptions → FlinkKafkaConsumerBase → AbstractRichFunction → RichFunction
  33. def open(arg0: OpenContext): Unit
    Definition Classes
    RichFunction
    Annotations
    @throws(classOf[java.lang.Exception])
  34. def run(arg0: SourceContext[T]): Unit
    Definition Classes
    FlinkKafkaConsumerBase → SourceFunction
    Annotations
    @throws(classOf[java.lang.Exception])
  35. def setCommitOffsetsOnCheckpoints(arg0: Boolean): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  36. def setRuntimeContext(arg0: RuntimeContext): Unit
    Definition Classes
    AbstractRichFunction → RichFunction
  37. def setStartFromEarliest(): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  38. def setStartFromGroupOffsets(): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  39. def setStartFromLatest(): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  40. def setStartFromSpecificOffsets(arg0: Map[KafkaTopicPartition, Long]): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  41. def setStartFromTimestamp(arg0: Long): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
  42. final def snapshotState(arg0: FunctionSnapshotContext): Unit
    Definition Classes
    FlinkKafkaConsumerBase → CheckpointedFunction
    Annotations
    @throws(classOf[java.lang.Exception])
  43. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  44. def toString(): String
    Definition Classes
    AnyRef → Any
  45. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  46. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException]) @native()
  47. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])

Deprecated Value Members

  1. def assignTimestampsAndWatermarks(arg0: AssignerWithPeriodicWatermarks[T]): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
    Annotations
    @Deprecated
    Deprecated
  2. def assignTimestampsAndWatermarks(arg0: AssignerWithPunctuatedWatermarks[T]): FlinkKafkaConsumerBase[T]
    Definition Classes
    FlinkKafkaConsumerBase
    Annotations
    @Deprecated
    Deprecated
  3. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.Throwable]) @Deprecated
    Deprecated

    (Since version 9)

Inherited from LazyLogging

Inherited from FlinkKafkaConsumer[T]

Inherited from FlinkKafkaConsumerBase[T]

Inherited from CheckpointedFunction

Inherited from ResultTypeQueryable[T]

Inherited from CheckpointListener

Inherited from RichParallelSourceFunction[T]

Inherited from ParallelSourceFunction[T]

Inherited from SourceFunction[T]

Inherited from AbstractRichFunction

Inherited from RichFunction

Inherited from Function

Inherited from Serializable

Inherited from AnyRef

Inherited from Any

Ungrouped