Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-19 Thread Shushant Arora
DefaultDecoder, >>>>>>>>>> kafka.serializer.DefaultDecoder, byte[][]>{ >>>>>>>>>> >>>>>>>>>> public CustomDirectKafkaInputDstream( >>>>>>>>>> StreamingContext ssc_, >>>

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-19 Thread Cody Koeninger
>>>>>>>>> StreamingContext ssc_, >>>>>>>>> Map kafkaParams, >>>>>>>>> Map fromOffsets, >>>>>>>>> Function1, byte[][]> >>>>>>>>> messageHandler, >>&

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-19 Thread Shushant Arora
nce$1, >>>>>>>> evidence$2, >>>>>>>> evidence$3, evidence$4, evidence$5); >>>>>>>> } >>>>>>>> @Override >>>>>>>> public Option>>>>>>> DefaultDecoder, byte[][]>

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Shushant Arora
r, byte[][]>> compute( >>>>>>> Time validTime) { >>>>>>> int processe=processedCounter.value(); >>>>>>> int failed = failedProcessingsCounter.value(); >>>>>>> if((processed==failed)){ >>>>&g

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Cody Koeninger
>>>>>> return super.compute(validTime); >>>>>> } >>>>>> } >>>>>> >>>>>> >>>>>> >>>>>> To create this stream >>>>>> I am using >>>>>> scala.collecti

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Shushant Arora
;>>>> >>>>> To create this stream >>>>> I am using >>>>> scala.collection.immutable.Map scalakafkaParams = >>>>> JavaConverters.mapAsScalaMapConverter(kafkaParams).asScala().toMap(Predef.>>>> String>>conforms());

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Shushant Arora
t; new Function, byte[][]>() { >>>> ..}); >>>> JavaDStream directKafkaStream = new >>>> CustomDirectKafkaInputDstream(jssc,scalakafkaParams ,scalaktopicOffsetMap, >>>> handler,byte[].class,byte[].class, kafka.serializer.DefaultDecoder.clas

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Cody Koeninger
afka.serializer.DefaultDecoder.class, >>> kafka.serializer.DefaultDecoder.class,byte[][].class); >>> >>> >>> >>> How to pass classTag to constructor in CustomDirectKafkaInputDstream ? >>> And how to use Function instead of Function1 ? >>> >>> >

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-18 Thread Shushant Arora
but you could create your own >>> subclass of the DStream that returns None for compute() under certain >>> conditions. >>> >>> >>> >>> On Wed, Aug 12, 2015 at 1:03 PM, Shushant Arora < >>> shushantaror...@gmail.com> wrote: >>&

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-17 Thread Cody Koeninger
ushantaror...@gmail.com> wrote: >> >>> Hi Cody >>> >>> Can you help here if streaming 1.3 has any api for not consuming any >>> message in next few runs? >>> >>> Thanks >>> >>> -- Forwarded message -

Re: spark streaming 1.3 doubts(force it to not consume anything)

2015-08-17 Thread Shushant Arora
gt;> Hi Cody >> >> Can you help here if streaming 1.3 has any api for not consuming any >> message in next few runs? >> >> Thanks >> >> ------ Forwarded message ---------- >> From: Shushant Arora >> Date: Wed, Aug 12, 2015 at 11:23 PM >&