public class KafkaUtilsPythonHelper
extends Object
The zero-arg constructor helps instantiate this class from the Class object classOf[KafkaUtilsPythonHelper].newInstance(), and the createStream() takes care of known parameters instead of passing them from Python
Constructor and Description |
---|
KafkaUtilsPythonHelper() |
Modifier and Type | Method and Description |
---|---|
Broker |
createBroker(String host,
Integer port) |
JavaPairInputDStream<byte[],byte[]> |
createDirectStream(JavaStreamingContext jssc,
java.util.Map<String,String> kafkaParams,
java.util.Set<String> topics,
java.util.Map<kafka.common.TopicAndPartition,Long> fromOffsets) |
OffsetRange |
createOffsetRange(String topic,
Integer partition,
Long fromOffset,
Long untilOffset) |
JavaPairRDD<byte[],byte[]> |
createRDD(JavaSparkContext jsc,
java.util.Map<String,String> kafkaParams,
java.util.List<OffsetRange> offsetRanges,
java.util.Map<kafka.common.TopicAndPartition,Broker> leaders) |
JavaPairReceiverInputDStream<byte[],byte[]> |
createStream(JavaStreamingContext jssc,
java.util.Map<String,String> kafkaParams,
java.util.Map<String,Integer> topics,
StorageLevel storageLevel) |
kafka.common.TopicAndPartition |
createTopicAndPartition(String topic,
Integer partition) |
public JavaPairReceiverInputDStream<byte[],byte[]> createStream(JavaStreamingContext jssc, java.util.Map<String,String> kafkaParams, java.util.Map<String,Integer> topics, StorageLevel storageLevel)
public JavaPairRDD<byte[],byte[]> createRDD(JavaSparkContext jsc, java.util.Map<String,String> kafkaParams, java.util.List<OffsetRange> offsetRanges, java.util.Map<kafka.common.TopicAndPartition,Broker> leaders)
public JavaPairInputDStream<byte[],byte[]> createDirectStream(JavaStreamingContext jssc, java.util.Map<String,String> kafkaParams, java.util.Set<String> topics, java.util.Map<kafka.common.TopicAndPartition,Long> fromOffsets)
public OffsetRange createOffsetRange(String topic, Integer partition, Long fromOffset, Long untilOffset)
public kafka.common.TopicAndPartition createTopicAndPartition(String topic, Integer partition)
public Broker createBroker(String host, Integer port)