Packages

case class KafkaInput(name: String, brokers: String, topic: Option[String], topicPattern: Option[String], consumerGroup: Option[String], options: Option[Map[String, String]], schemaRegistryUrl: Option[String], schemaSubject: Option[String], schemaId: Option[Int], isBatch: Option[Boolean], format: Option[String], incremental: Option[Incremental]) extends IncrementalReader with DatasourceReader with Product with Serializable

Represents a Kafka data source input configuration.

name

The name of the Kafka input.

brokers

The Kafka brokers to connect to.

topic

Optional Kafka topic to subscribe to.

topicPattern

Optional Kafka topic pattern to subscribe to.

consumerGroup

Optional Kafka consumer group ID.

options

Optional additional Kafka options.

schemaRegistryUrl

Optional URL for the Avro schema registry.

schemaSubject

Optional subject name for the Avro schema.

schemaId

Optional ID for the Avro schema.

isBatch

Optional flag indicating batch processing.

format

Optional data format for Kafka messages.

incremental

Optional incremental configuration.

Linear Supertypes
Serializable, Serializable, Product, Equals, DatasourceReader, Reader, IncrementalReader, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. KafkaInput
  2. Serializable
  3. Serializable
  4. Product
  5. Equals
  6. DatasourceReader
  7. Reader
  8. IncrementalReader
  9. AnyRef
  10. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. All

Instance Constructors

  1. new KafkaInput(name: String, brokers: String, topic: Option[String], topicPattern: Option[String], consumerGroup: Option[String], options: Option[Map[String, String]], schemaRegistryUrl: Option[String], schemaSubject: Option[String], schemaId: Option[Int], isBatch: Option[Boolean], format: Option[String], incremental: Option[Incremental])

    name

    The name of the Kafka input.

    brokers

    The Kafka brokers to connect to.

    topic

    Optional Kafka topic to subscribe to.

    topicPattern

    Optional Kafka topic pattern to subscribe to.

    consumerGroup

    Optional Kafka consumer group ID.

    options

    Optional additional Kafka options.

    schemaRegistryUrl

    Optional URL for the Avro schema registry.

    schemaSubject

    Optional subject name for the Avro schema.

    schemaId

    Optional ID for the Avro schema.

    isBatch

    Optional flag indicating batch processing.

    format

    Optional data format for Kafka messages.

    incremental

    Optional incremental configuration.

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. val brokers: String
  6. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native()
  7. val consumerGroup: Option[String]
  8. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  9. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( classOf[java.lang.Throwable] )
  10. val format: Option[String]
  11. val fromConfluentAvro: Boolean
  12. final def getClass(): Class[_]
    Definition Classes
    AnyRef → Any
    Annotations
    @native()
  13. val incremental: Option[Incremental]
  14. val isBatch: Option[Boolean]
  15. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  16. lazy val log: Logger
    Annotations
    @transient()
  17. val name: String

    The name of the data reader.

    The name of the data reader.

    Definition Classes
    KafkaInputReader
  18. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  19. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native()
  20. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native()
  21. val options: Option[Map[String, String]]
  22. def persistIncrementalState(): Unit

    Persists the incremental state if any.

    Persists the incremental state if any.

    Definition Classes
    IncrementalReader
  23. def read(sparkSession: SparkSession): DataFrame

    Reads data from the Kafka data source and returns a DataFrame.

    Reads data from the Kafka data source and returns a DataFrame.

    sparkSession

    The SparkSession instance.

    returns

    The DataFrame containing the data from the Kafka data source.

    Definition Classes
    KafkaInputReader
  24. def readIncremental(df: DataFrame, incremental: Option[Incremental]): DataFrame

    Reads data from a DataFrame with optional incremental settings.

    Reads data from a DataFrame with optional incremental settings.

    df

    The DataFrame to read data from.

    incremental

    Optional Incremental settings to apply.

    returns

    A new DataFrame after applying incremental settings if provided, otherwise the original DataFrame.

    Definition Classes
    IncrementalReader
  25. val schemaId: Option[Int]
  26. val schemaRegistryUrl: Option[String]
  27. val schemaSubject: Option[String]
  28. val subject: String
  29. final def synchronized[T0](arg0: ⇒ T0): T0
    Definition Classes
    AnyRef
  30. val topic: Option[String]
  31. val topicPattern: Option[String]
  32. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  33. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  34. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native()

Inherited from Serializable

Inherited from Serializable

Inherited from Product

Inherited from Equals

Inherited from DatasourceReader

Inherited from Reader

Inherited from IncrementalReader

Inherited from AnyRef

Inherited from Any

Ungrouped