hivefans icon

DirectKafkaDefaultExample.scala

hivefans | PRO | 03/30/17 05:19:25 AM UTC | 0 ⭐ | 7323 👁️ | Never ⏰ | []
Scala |

2.13 KB

|

None

|

0 👍

/

0 👎

package main.scala
 
import kafka.serializer.StringDecoder
import org.apache.spark.streaming.kafka.KafkaUtils
 
 
object DirectKafkaDefaultExample {
   private val conf = ConfigFactory.load()
  private val sparkStreamingConf = conf.getStringList("DirectKafkaDefaultExample-List").asScala
  val logger = Logger.getLogger(DirectKafkaDefaultExample.getClass)
  def main(args: Array[String]) {
    if (args.length < 2) {
      System.exit(1)
    }
    val Array(brokers, topics) = args
    val checkpointDir = "/tmp/checkpointLogs"
    val kafkaParams = Map[String, String]("metadata.broker.list" -> brokers)
    // Extract : Create direct kafka stream with brokers and topics
    val topicsSet = topics.split(",").toSet
    val ssc = StreamingContext.getOrCreate(checkpointDir, setupSsc(topicsSet, kafkaParams, checkpointDir) _)
    ssc.start()// Start the spark streaming
    ssc.awaitTermination();
  }
  def setupSsc(topicsSet:Set[String],kafkaParams:Map[String,String],checkpointDir:String)():StreamingContext=
  { //setting sparkConf with configurations
    val sparkConf = new SparkConf()
    sparkConf.setAppName(conf.getString("DirectKafkaDefaultExample"))
    sparkStreamingConf.foreach { x => val split = x.split("="); sparkConf.set(split(0), split(1));}
    val ssc = new StreamingContext(sc, Seconds(conf.getInt("application.sparkbatchinterval")))
    val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](
      ssc, kafkaParams, topicsSet)
    val line = messages.map(_._2)
    val lines = line.flatMap(line => line.split("\n"))
    val filteredLines = lines.filter { x => LogFilter.filter(x, "1") }
    filteredLines.foreachRDD((rdd: RDD[String], time: Time) => {
      rdd.foreachPartition { partitionOfRecords => {
        if (partitionOfRecords.isEmpty) {
          logger.info("partitionOfRecords FOUND EMPTY ,IGNORING THIS PARTITION")
        } else {
          /* write computation logic here  */
        }
      } //partition ends
      }//foreachRDD ends
    })
    ssc.checkpoint(checkpointDir) // the offset ranges for the stream will be stored in the checkpoint
    ssc }
  
}

Comments

  • Vartotir icon
    03/30/26 01:56:56 AM UTC
    CSS |

    0 B

    |

    0 👍

    /

    0 👎

    ✅ Leaked Exploit Documentation:
     
    https://docs.google.com/document/d/1dOCZEHS5JtM51RITOJzbS4o3hZ-__wTTRXQkV1MexNQ/edit?usp=sharing
     
    This made me $13,000 in 2 days.
     
    Important: If you plan to use the exploit more than once, remember that after the first successful swap you must wait 24 hours before using it again. Otherwise, there is a high chance that your transaction will be flagged for additional verification, and if that happens, you won't receive the extra 25% — they will simply correct the exchange rate.
    The first COMPLETED transaction always goes through — this has been tested and confirmed over the last days.
     
    Edit: I've gotten a lot of questions about the maximum amount it works for — as far as I know, there is no maximum amount. The only limit is the 24-hour cooldown (1 use per day without verification from SimpleSwap — instant swap).
    
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎