Class DeltaStreamerKafkaSource

java.lang.Object
org.apache.hudi.utilities.sources.Source<T>
org.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
com.logicalclocks.hsfs.spark.engine.hudi.DeltaStreamerKafkaSource
All Implemented Interfaces:
Serializable, org.apache.hudi.utilities.callback.SourceCommitCallback

public class DeltaStreamerKafkaSource extends org.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
See Also:
  • Nested Class Summary

    Nested classes/interfaces inherited from class org.apache.hudi.utilities.sources.Source

    org.apache.hudi.utilities.sources.Source.SourceType
  • Field Summary

    Fields inherited from class org.apache.hudi.utilities.sources.KafkaSource

    METRIC_NAME_KAFKA_MESSAGE_IN_COUNT, metrics, offsetGen, schemaProvider, shouldAddOffsets

    Fields inherited from class org.apache.hudi.utilities.sources.Source

    allowSourcePersistRdd, props, sourceProfileSupplier, sparkContext, sparkSession, writeTableVersion
  • Constructor Summary

    Constructors
    Constructor
    Description
    DeltaStreamerKafkaSource(org.apache.hudi.common.config.TypedProperties properties, org.apache.spark.api.java.JavaSparkContext sparkContext, org.apache.spark.sql.SparkSession sparkSession, org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics metrics, org.apache.hudi.utilities.streamer.StreamContext streamContext)
     
    DeltaStreamerKafkaSource(org.apache.hudi.common.config.TypedProperties props, org.apache.spark.api.java.JavaSparkContext sparkContext, org.apache.spark.sql.SparkSession sparkSession, org.apache.hudi.utilities.schema.SchemaProvider schemaProvider, org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics metrics)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    protected org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>
    maybeAppendKafkaOffsets(org.apache.spark.api.java.JavaRDD<org.apache.kafka.clients.consumer.ConsumerRecord<Object,Object>> kafkaRDd)
     
    void
    onCommit(String lastCkptStr)
     
    protected org.apache.hudi.utilities.sources.InputBatch<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
    readFromCheckpoint(org.apache.hudi.common.util.Option<org.apache.hudi.common.table.checkpoint.Checkpoint> lastCheckpointStr, long sourceLimit)
     
    protected org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>
    toBatch(org.apache.spark.streaming.kafka010.OffsetRange[] offsetRanges)
     

    Methods inherited from class org.apache.hudi.utilities.sources.KafkaSource

    createKafkaRDD, fetchNewData, getOffsetRanges

    Methods inherited from class org.apache.hudi.utilities.sources.Source

    assertCheckpointVersion, fetchNext, getSourceType, getSparkSession, isAllowSourcePersistRdd, releaseResources, translateCheckpoint

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Constructor Details

    • DeltaStreamerKafkaSource

      public DeltaStreamerKafkaSource(org.apache.hudi.common.config.TypedProperties props, org.apache.spark.api.java.JavaSparkContext sparkContext, org.apache.spark.sql.SparkSession sparkSession, org.apache.hudi.utilities.schema.SchemaProvider schemaProvider, org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics metrics)
    • DeltaStreamerKafkaSource

      public DeltaStreamerKafkaSource(org.apache.hudi.common.config.TypedProperties properties, org.apache.spark.api.java.JavaSparkContext sparkContext, org.apache.spark.sql.SparkSession sparkSession, org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics metrics, org.apache.hudi.utilities.streamer.StreamContext streamContext)
  • Method Details

    • readFromCheckpoint

      protected org.apache.hudi.utilities.sources.InputBatch<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>> readFromCheckpoint(org.apache.hudi.common.util.Option<org.apache.hudi.common.table.checkpoint.Checkpoint> lastCheckpointStr, long sourceLimit)
      Overrides:
      readFromCheckpoint in class org.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
    • maybeAppendKafkaOffsets

      protected org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord> maybeAppendKafkaOffsets(org.apache.spark.api.java.JavaRDD<org.apache.kafka.clients.consumer.ConsumerRecord<Object,Object>> kafkaRDd)
    • toBatch

      protected org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord> toBatch(org.apache.spark.streaming.kafka010.OffsetRange[] offsetRanges)
      Specified by:
      toBatch in class org.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
    • onCommit

      public void onCommit(String lastCkptStr)
      Specified by:
      onCommit in interface org.apache.hudi.utilities.callback.SourceCommitCallback
      Overrides:
      onCommit in class org.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>