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, shouldAddOffsetsFields inherited from class org.apache.hudi.utilities.sources.Source
allowSourcePersistRdd, props, sourceProfileSupplier, sparkContext, sparkSession, writeTableVersion -
Constructor Summary
ConstructorsConstructorDescriptionDeltaStreamerKafkaSource(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 TypeMethodDescriptionprotected 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) voidprotected 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, getOffsetRangesMethods inherited from class org.apache.hudi.utilities.sources.Source
assertCheckpointVersion, fetchNext, getSourceType, getSparkSession, isAllowSourcePersistRdd, releaseResources, translateCheckpoint
-
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:
readFromCheckpointin classorg.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
-
maybeAppendKafkaOffsets
-
toBatch
protected org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord> toBatch(org.apache.spark.streaming.kafka010.OffsetRange[] offsetRanges) - Specified by:
toBatchin classorg.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
-
onCommit
- Specified by:
onCommitin interfaceorg.apache.hudi.utilities.callback.SourceCommitCallback- Overrides:
onCommitin classorg.apache.hudi.utilities.sources.KafkaSource<org.apache.spark.api.java.JavaRDD<org.apache.avro.generic.GenericRecord>>
-