How to use specific ElasticSearch query according to SparkStreaming-Kafka messageElasticsearch query to return all recordselasticsearch bool query combine must with ORSparkStreaming keep processing even no data in kafkaIs it possible to obtain specific message offset in Kafka+SparkStreaming?SparkStreaming/Kafka Offset handlingSparkStreaming put a kafka message into Hbasedo not want string as type when using foreach in scala spark streaming?[SparkStreaming]Kafka ConsumerRecord is not serializableSparkStreaming & Kafka: value reduceByKey is not a member of org.apache.spark.streaming.dstream.DStream[Any]Splitting Kafka Message Line by line in Spark Structured Streaming

How can the US president give an order to a civilian?

Why is Skinner so awkward in Hot Fuzz?

Is there a risk to write an invitation letter for a stranger to obtain a Czech (Schengen) visa?

How to write a nice frame challenge?

How "fast" does astronomical events happen?

Are their examples of rowers who also fought?

How can I detect if I'm in a subshell?

How to know whether to write accidentals as sharps or flats?

Is there any effect in D&D 5e that cannot be undone?

How could I create a situation in which a PC has to make a saving throw or be forced to pet a dog?

Leveraging cash for buying car

Explicit direct #include vs. Non-contractual transitive #include

How can I maintain game balance while allowing my player to craft genuinely useful items?

My student in one course asks for paid tutoring in another course. Appropriate?

How Linux command "mount -a" works

How can a flywheel makes engine runs smoothly?

How to make a villain when your PCs are villains?

Right indicator flash-frequency has increased and rear-right bulb is out

How to ask if I can mow my neighbor's lawn

How do credit card companies know what type of business I'm paying for?

Testing thermite for chemical properties

How to prevent cables getting intertwined

How to avoid offending original culture when making conculture inspired from original

Why can't I craft scaffolding in Minecraft 1.14?



How to use specific ElasticSearch query according to SparkStreaming-Kafka message


Elasticsearch query to return all recordselasticsearch bool query combine must with ORSparkStreaming keep processing even no data in kafkaIs it possible to obtain specific message offset in Kafka+SparkStreaming?SparkStreaming/Kafka Offset handlingSparkStreaming put a kafka message into Hbasedo not want string as type when using foreach in scala spark streaming?[SparkStreaming]Kafka ConsumerRecord is not serializableSparkStreaming & Kafka: value reduceByKey is not a member of org.apache.spark.streaming.dstream.DStream[Any]Splitting Kafka Message Line by line in Spark Structured Streaming






.everyoneloves__top-leaderboard:empty,.everyoneloves__mid-leaderboard:empty,.everyoneloves__bot-mid-leaderboard:empty height:90px;width:728px;box-sizing:border-box;








0















I'm using SparkStreaming-Kafka, and want to support ElasticSearch real_time search with specific query, which is from Kafka message.

Code like bellow:



def creatingFuncTest():StreamingContext=

val ssc = new StreamingContext(sc, Seconds(duration.toInt))
ssc.checkpoint(checkpointDir)
val kafkaParams = KafkaUtil.getKafkaParam(brokers, appName)

val topics = actionTopicList.split(",").toSet

val foodMessages = KafkaUtils
.createDirectStream[String, String, StringDecoder, StringDecoder](
ssc,
kafkaParams,
topics
)


val foodBatch: DStream[(String,Float, Float)] =
foodMessages
.filter(_._2.nonEmpty)
.map msg =>
try
println("___________ msg :" + msg._2)
val gson = new Gson()
val vo = gson.fromJson(msg._2, classOf[PoiMsg])
(vo.person_id.toString, vo.latitude.toFloat,vo.longitude.toFloat)
catch
case e: Exception =>
println("____________" + e.getMessage)
("", 0.0f,0.0f)


.filter(_._1.nonEmpty)



foodBatch.foreachRDD(row =>
row.foreach(t =>
var lat = t._2
var lon = t._3


val query:String =s"""
"filter" :
"geo_distance" : //...
"distance" : "200km",
"pin.location" : "lat" : "$lat", "lon" : "#lon"

"""

val rdd = sc.esRDD("recommend_diet_menu/fooddocument", query)

println(rdd.count())

)
)
ssc




I know in RDD, it's wrong to generate new RDD, but what's the right way?










share|improve this question






























    0















    I'm using SparkStreaming-Kafka, and want to support ElasticSearch real_time search with specific query, which is from Kafka message.

    Code like bellow:



    def creatingFuncTest():StreamingContext=

    val ssc = new StreamingContext(sc, Seconds(duration.toInt))
    ssc.checkpoint(checkpointDir)
    val kafkaParams = KafkaUtil.getKafkaParam(brokers, appName)

    val topics = actionTopicList.split(",").toSet

    val foodMessages = KafkaUtils
    .createDirectStream[String, String, StringDecoder, StringDecoder](
    ssc,
    kafkaParams,
    topics
    )


    val foodBatch: DStream[(String,Float, Float)] =
    foodMessages
    .filter(_._2.nonEmpty)
    .map msg =>
    try
    println("___________ msg :" + msg._2)
    val gson = new Gson()
    val vo = gson.fromJson(msg._2, classOf[PoiMsg])
    (vo.person_id.toString, vo.latitude.toFloat,vo.longitude.toFloat)
    catch
    case e: Exception =>
    println("____________" + e.getMessage)
    ("", 0.0f,0.0f)


    .filter(_._1.nonEmpty)



    foodBatch.foreachRDD(row =>
    row.foreach(t =>
    var lat = t._2
    var lon = t._3


    val query:String =s"""
    "filter" :
    "geo_distance" : //...
    "distance" : "200km",
    "pin.location" : "lat" : "$lat", "lon" : "#lon"

    """

    val rdd = sc.esRDD("recommend_diet_menu/fooddocument", query)

    println(rdd.count())

    )
    )
    ssc




    I know in RDD, it's wrong to generate new RDD, but what's the right way?










    share|improve this question


























      0












      0








      0








      I'm using SparkStreaming-Kafka, and want to support ElasticSearch real_time search with specific query, which is from Kafka message.

      Code like bellow:



      def creatingFuncTest():StreamingContext=

      val ssc = new StreamingContext(sc, Seconds(duration.toInt))
      ssc.checkpoint(checkpointDir)
      val kafkaParams = KafkaUtil.getKafkaParam(brokers, appName)

      val topics = actionTopicList.split(",").toSet

      val foodMessages = KafkaUtils
      .createDirectStream[String, String, StringDecoder, StringDecoder](
      ssc,
      kafkaParams,
      topics
      )


      val foodBatch: DStream[(String,Float, Float)] =
      foodMessages
      .filter(_._2.nonEmpty)
      .map msg =>
      try
      println("___________ msg :" + msg._2)
      val gson = new Gson()
      val vo = gson.fromJson(msg._2, classOf[PoiMsg])
      (vo.person_id.toString, vo.latitude.toFloat,vo.longitude.toFloat)
      catch
      case e: Exception =>
      println("____________" + e.getMessage)
      ("", 0.0f,0.0f)


      .filter(_._1.nonEmpty)



      foodBatch.foreachRDD(row =>
      row.foreach(t =>
      var lat = t._2
      var lon = t._3


      val query:String =s"""
      "filter" :
      "geo_distance" : //...
      "distance" : "200km",
      "pin.location" : "lat" : "$lat", "lon" : "#lon"

      """

      val rdd = sc.esRDD("recommend_diet_menu/fooddocument", query)

      println(rdd.count())

      )
      )
      ssc




      I know in RDD, it's wrong to generate new RDD, but what's the right way?










      share|improve this question
















      I'm using SparkStreaming-Kafka, and want to support ElasticSearch real_time search with specific query, which is from Kafka message.

      Code like bellow:



      def creatingFuncTest():StreamingContext=

      val ssc = new StreamingContext(sc, Seconds(duration.toInt))
      ssc.checkpoint(checkpointDir)
      val kafkaParams = KafkaUtil.getKafkaParam(brokers, appName)

      val topics = actionTopicList.split(",").toSet

      val foodMessages = KafkaUtils
      .createDirectStream[String, String, StringDecoder, StringDecoder](
      ssc,
      kafkaParams,
      topics
      )


      val foodBatch: DStream[(String,Float, Float)] =
      foodMessages
      .filter(_._2.nonEmpty)
      .map msg =>
      try
      println("___________ msg :" + msg._2)
      val gson = new Gson()
      val vo = gson.fromJson(msg._2, classOf[PoiMsg])
      (vo.person_id.toString, vo.latitude.toFloat,vo.longitude.toFloat)
      catch
      case e: Exception =>
      println("____________" + e.getMessage)
      ("", 0.0f,0.0f)


      .filter(_._1.nonEmpty)



      foodBatch.foreachRDD(row =>
      row.foreach(t =>
      var lat = t._2
      var lon = t._3


      val query:String =s"""
      "filter" :
      "geo_distance" : //...
      "distance" : "200km",
      "pin.location" : "lat" : "$lat", "lon" : "#lon"

      """

      val rdd = sc.esRDD("recommend_diet_menu/fooddocument", query)

      println(rdd.count())

      )
      )
      ssc




      I know in RDD, it's wrong to generate new RDD, but what's the right way?







      apache-spark elasticsearch spark-streaming spark-streaming-kafka






      share|improve this question















      share|improve this question













      share|improve this question




      share|improve this question








      edited Mar 25 at 6:39









      pirho

      5,299122132




      5,299122132










      asked Mar 25 at 3:59









      laclac

      3016




      3016






















          0






          active

          oldest

          votes












          Your Answer






          StackExchange.ifUsing("editor", function ()
          StackExchange.using("externalEditor", function ()
          StackExchange.using("snippets", function ()
          StackExchange.snippets.init();
          );
          );
          , "code-snippets");

          StackExchange.ready(function()
          var channelOptions =
          tags: "".split(" "),
          id: "1"
          ;
          initTagRenderer("".split(" "), "".split(" "), channelOptions);

          StackExchange.using("externalEditor", function()
          // Have to fire editor after snippets, if snippets enabled
          if (StackExchange.settings.snippets.snippetsEnabled)
          StackExchange.using("snippets", function()
          createEditor();
          );

          else
          createEditor();

          );

          function createEditor()
          StackExchange.prepareEditor(
          heartbeatType: 'answer',
          autoActivateHeartbeat: false,
          convertImagesToLinks: true,
          noModals: true,
          showLowRepImageUploadWarning: true,
          reputationToPostImages: 10,
          bindNavPrevention: true,
          postfix: "",
          imageUploader:
          brandingHtml: "Powered by u003ca class="icon-imgur-white" href="https://imgur.com/"u003eu003c/au003e",
          contentPolicyHtml: "User contributions licensed under u003ca href="https://creativecommons.org/licenses/by-sa/3.0/"u003ecc by-sa 3.0 with attribution requiredu003c/au003e u003ca href="https://stackoverflow.com/legal/content-policy"u003e(content policy)u003c/au003e",
          allowUrls: true
          ,
          onDemand: true,
          discardSelector: ".discard-answer"
          ,immediatelyShowMarkdownHelp:true
          );



          );













          draft saved

          draft discarded


















          StackExchange.ready(
          function ()
          StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f55331094%2fhow-to-use-specific-elasticsearch-query-according-to-sparkstreaming-kafka-messag%23new-answer', 'question_page');

          );

          Post as a guest















          Required, but never shown

























          0






          active

          oldest

          votes








          0






          active

          oldest

          votes









          active

          oldest

          votes






          active

          oldest

          votes















          draft saved

          draft discarded
















































          Thanks for contributing an answer to Stack Overflow!


          • Please be sure to answer the question. Provide details and share your research!

          But avoid


          • Asking for help, clarification, or responding to other answers.

          • Making statements based on opinion; back them up with references or personal experience.

          To learn more, see our tips on writing great answers.




          draft saved


          draft discarded














          StackExchange.ready(
          function ()
          StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f55331094%2fhow-to-use-specific-elasticsearch-query-according-to-sparkstreaming-kafka-messag%23new-answer', 'question_page');

          );

          Post as a guest















          Required, but never shown





















































          Required, but never shown














          Required, but never shown












          Required, but never shown







          Required, but never shown

































          Required, but never shown














          Required, but never shown












          Required, but never shown







          Required, but never shown







          Popular posts from this blog

          Kamusi Yaliyomo Aina za kamusi | Muundo wa kamusi | Faida za kamusi | Dhima ya picha katika kamusi | Marejeo | Tazama pia | Viungo vya nje | UrambazajiKuhusu kamusiGo-SwahiliWiki-KamusiKamusi ya Kiswahili na Kiingerezakuihariri na kuongeza habari

          SQL error code 1064 with creating Laravel foreign keysForeign key constraints: When to use ON UPDATE and ON DELETEDropping column with foreign key Laravel error: General error: 1025 Error on renameLaravel SQL Can't create tableLaravel Migration foreign key errorLaravel php artisan migrate:refresh giving a syntax errorSQLSTATE[42S01]: Base table or view already exists or Base table or view already exists: 1050 Tableerror in migrating laravel file to xampp serverSyntax error or access violation: 1064:syntax to use near 'unsigned not null, modelName varchar(191) not null, title varchar(191) not nLaravel cannot create new table field in mysqlLaravel 5.7:Last migration creates table but is not registered in the migration table

          은진 송씨 목차 역사 본관 분파 인물 조선 왕실과의 인척 관계 집성촌 항렬자 인구 같이 보기 각주 둘러보기 메뉴은진 송씨세종실록 149권, 지리지 충청도 공주목 은진현