Error While Inserting Data into Cassandra -


i learning apache kafka-storm-cassandra integration. reading json string kafka cluster using kafka spout.then passing bolt parses json , emits needed value 2nd bolt writes cassandra db.

but getting these errors.

java.lang.runtimeexception: com.datastax.driver.core.exceptions.syntaxerror: line 1:154 mismatched character '' expecting '''

java.lang.runtimeexception: com.datastax.driver.core.exceptions.syntaxerror: line 1:154 mismatched character '<eof>' expecting ''' @  backtype.storm.utils.disruptorqueue.consumebatchtocursor(disruptorqueue.java:128) @  backtype.storm.utils.disruptorqueue.consumebatchwhenavailable(disruptorqueue.java:99) @  backtype.storm.disruptor$consume_batch_when_available.invoke(disruptor.clj:80) @  backtype.storm.daemon.executor$fn__4722$fn__4734$fn__4781.invoke(executor.clj:748) @ backtype.storm.util$async_loop$fn__458.invoke(util.clj:463) @ clojure.lang.afn.run(afn.java:24) @ java.lang.thread.run(thread.java:745)  caused by: com.datastax.driver.core.exceptions.syntaxerror: line 1:154 mismatched character '<eof>' expecting ''' @ com.datastax.driver.core.exceptions.syntaxerror.copy(syntaxerror.java:35) @ com.datastax.driver.core.defaultresultsetfuture.extractcausefromexecutionexception(defaultresultsetfuture.java:289) @  com.datastax.driver.core.defaultresultsetfuture.getuninterruptibly(defaultresultsetfuture.java:205) @ com.datastax.driver.core.abstractsession.execute(abstractsession.java:52) @ com.datastax.driver.core.abstractsession.execute(abstractsession.java:36) @ bolts.wordcounter.execute(wordcounter.java:103) @ backtype.storm.topology.basicboltexecutor.execute(basicboltexecutor.java:50) @  backtype.storm.daemon.executor$fn__4722$tuple_action_fn__4724.invoke(executor.clj:633) @ backtype.storm.daemon.executor$mk_task_receiver$fn__4645.invoke(executor.clj:401) @ backtype.storm.disruptor$clojure_handler$reify__1446.onevent(disruptor.clj:58) @  backtype.storm.utils.disruptorqueue.consumebatchtocursor(disruptorqueue.java:120) ... 6 more caused by: com.datastax.driver.core.exceptions.syntaxerror: line 1:154 mismatched character '<eof>' expecting ''' @ com.datastax.driver.core.responses$error.asexception(responses.java:101) @ com.datastax.driver.core.defaultresultsetfuture.onset(defaultresultsetfuture.java:140) @  com.datastax.driver.core.requesthandler.setfinalresult(requesthandler.java:293) @ com.datastax.driver.core.requesthandler.onset(requesthandler.java:455) @  com.datastax.driver.core.connection$dispatcher.messagereceived(connection.java:734) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @ org.jboss.netty.handler.timeout.idlestateawarechannelupstreamhandler.handleupstream(idlestateawarechannelupstreamhandler.java:36) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @  org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.handler.timeout.idlestatehandler.messagereceived(idlestatehandler.java:294) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @  org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:296) @ org.jboss.netty.handler.codec.oneone.onetoonedecoder.handleupstream(onetoonedecoder.java:70) @  org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:296) @ org.jboss.netty.handler.codec.frame.framedecoder.unfoldandfiremessagereceived(framedecoder.java:462) @  org.jboss.netty.handler.codec.frame.framedecoder.calldecode(framedecoder.java:443) @ org.jboss.netty.handler.codec.frame.framedecoder.messagereceived(framedecoder.java:303) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:559) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:268) @  org.jboss.netty.channel.channels.firemessagereceived(channels.java:255) @ org.jboss.netty.channel.socket.nio.nioworker.read(nioworker.java:88) @ org.jboss.netty.channel.socket.nio.abstractnioworker.process(abstractnioworker.java:108) @ org.jboss.netty.channel.socket.nio.abstractnioselector.run(abstractnioselector.java:318) @ org.jboss.netty.channel.socket.nio.abstractnioworker.run(abstractnioworker.java:89) @ org.jboss.netty.channel.socket.nio.nioworker.run(nioworker.java:178) @  org.jboss.netty.util.threadrenamingrunnable.run(threadrenamingrunnable.java:108) @ org.jboss.netty.util.internal.deadlockproofworker$1.run(deadlockproofworker.java:42) @ java.util.concurrent.threadpoolexecutor.runworker(threadpoolexecutor.java:1145) @  java.util.concurrent.threadpoolexecutor$worker.run(threadpoolexecutor.java:615) ... 1 more 

tablealreadyexists error :

com.datastax.driver.core.exceptions.alreadyexistsexception: table query.productcount exists @ com.datastax.driver.core.exceptions.alreadyexistsexception.copy(alreadyexistsexception.java:85)

com.datastax.driver.core.exceptions.alreadyexistsexception: table query.productcount exists @ com.datastax.driver.core.exceptions.alreadyexistsexception.copy(alreadyexistsexception.java:85) @  com.datastax.driver.core.defaultresultsetfuture.extractcausefromexecutionexception(defaultresultsetfuture.java:289) @ com.datastax.driver.core.defaultresultsetfuture.getuninterruptibly(defaultresultsetfuture.java:205) @ com.datastax.driver.core.abstractsession.execute(abstractsession.java:52) @  com.datastax.driver.core.abstractsession.execute(abstractsession.java:36) @ bolts.wordcounter.prepare(wordcounter.java:77) @ backtype.storm.topology.basicboltexecutor.prepare(basicboltexecutor.java:43) @ backtype.storm.daemon.executor$fn__4722$fn__4734.invoke(executor.clj:692) @ backtype.storm.util$async_loop$fn__458.invoke(util.clj:461) @ clojure.lang.afn.run(afn.java:24) @ java.lang.thread.run(thread.java:745) caused by: com.datastax.driver.core.exceptions.alreadyexistsexception: table query.productcount exists @  com.datastax.driver.core.exceptions.alreadyexistsexception.copy(alreadyexistsexception.java:85) @  com.datastax.driver.core.responses$error.asexception(responses.java:105) @   com.datastax.driver.core.defaultresultsetfuture.onset(defaultresultsetfuture.java:140) @  com.datastax.driver.core.requesthandler.setfinalresult(requesthandler.java:293) @ com.datastax.driver.core.requesthandler.onset(requesthandler.java:455) @  com.datastax.driver.core.connection$dispatcher.messagereceived(connection.java:734) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @ org.jboss.netty.handler.timeout.idlestateawarechannelupstreamhandler.handleupstream(idlestateawarechannelupstreamhandler.java:36) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.handler.timeout.idlestatehandler.messagereceived(idlestatehandler.java:294) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:296) @ org.jboss.netty.handler.codec.oneone.onetoonedecoder.handleupstream(onetoonedecoder.java:70) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline$defaultchannelhandlercontext.sendupstream(defaultchannelpipeline.java:791) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:296) @ org.jboss.netty.handler.codec.frame.framedecoder.unfoldandfiremessagereceived(framedecoder.java:462) @ org.jboss.netty.handler.codec.frame.framedecoder.calldecode(framedecoder.java:443) @ org.jboss.netty.handler.codec.frame.framedecoder.messagereceived(framedecoder.java:303) @ org.jboss.netty.channel.simplechannelupstreamhandler.handleupstream(simplechannelupstreamhandler.java:70) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:564) @ org.jboss.netty.channel.defaultchannelpipeline.sendupstream(defaultchannelpipeline.java:559) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:268) @ org.jboss.netty.channel.channels.firemessagereceived(channels.java:255) @ org.jboss.netty.channel.socket.nio.nioworker.read(nioworker.java:88) @ org.jboss.netty.channel.socket.nio.abstractnioworker.process(abstractnioworker.java:108) @ org.jboss.netty.channel.socket.nio.abstractnioselector.run(abstractnioselector.java:318) @ org.jboss.netty.channel.socket.nio.abstractnioworker.run(abstractnioworker.java:89) @ org.jboss.netty.channel.socket.nio.nioworker.run(nioworker.java:178) @ org.jboss.netty.util.threadrenamingrunnable.run(threadrenamingrunnable.java:108) @ org.jboss.netty.util.internal.deadlockproofworker$1.run(deadlockproofworker.java:42) @ java.util.concurrent.threadpoolexecutor.runworker(threadpoolexecutor.java:1145) @ java.util.concurrent.threadpoolexecutor$worker.run(threadpoolexecutor.java:615) ... 1 more caused by: com.datastax.driver.core.exceptions.alreadyexistsexception: table query.productcount exists @ com.datastax.driver.core.responses$error$1.decode(responses.java:70) @ com.datastax.driver.core.responses$error$1.decode(responses.java:38) @ com.datastax.driver.core.message$protocoldecoder.decode(message.java:168) @ org.jboss.netty.handler.codec.oneone.onetoonedecoder.handleupstream(onetoonedecoder.java:66) ... 21 more 

my maintopology:

public class topologyquerycountermain {   static final logger logger = logger.getlogger(topologyquerycountermain.class);   private static final string spout_id = "querycounter";   public static void main(string[] args) throws alreadyaliveexception, invalidtopologyexception {      int numspoutexecutors = 1;     logger.debug("this spoutconfig");     kafkaspout kspout = querycounter();     topologybuilder builder = new topologybuilder();     logger.debug("this set spout");     builder.setspout(spout_id, kspout, numspoutexecutors);     logger.debug("this set bolt");     builder.setbolt("word-normalizer", new wordnormalizer())         .shufflegrouping(spout_id);     builder.setbolt("word-counter", new wordcounter(),1)         .shufflegrouping("word-normalizer", "stream1");       config conf = new config();     localcluster cluster = new localcluster();     logger.debug("this submit cluster");     conf.put(config.nimbus_host, "192.168.1.229");     conf.put(config.nimbus_thrift_port, 6627);      system.setproperty("storm.jar", "/home/ubuntu/workspace/querycounter/target/querycounter-0.0.1-snapshot.jar");     conf.setnumworkers(20);     conf.setmaxspoutpending(5000);      if (args != null && args.length > 0) {         stormsubmitter. submittopology(args[0], conf, builder.createtopology());     }      else     {            cluster.submittopology("querycounter", conf, builder.createtopology());         utils.sleep(10000);         cluster.killtopology("querycounter");         logger.debug("this shutdown cluster");         cluster.shutdown();     } }   private static kafkaspout querycounter() {     string zkhostport = "localhost:2181";     string topic = "randomquery";      string zkroot = "/querycounter";     string zkspoutid = "querycounter-spout";     zkhosts zkhosts = new zkhosts(zkhostport);      logger.debug("this inside kafka spout cluster");     spoutconfig spoutcfg = new spoutconfig(zkhosts, topic, zkroot, zkspoutid);     spoutcfg.scheme=new schemeasmultischeme(new stringscheme());     kafkaspout kafkaspout = new kafkaspout(spoutcfg);     return kafkaspout;   }  } 

normalizer bolt:

public class wordnormalizer extends basebasicbolt { static final logger logger = logger.getlogger(wordnormalizer.class); public void cleanup() {}  /**  * bolt receive line  * words file , process normalize line  *   * normalize put words in lower case  * , split line words in   */ public void execute(tuple input, basicoutputcollector collector) { string feed = input.getstring(0);      string searchterm = null;     string pageno = null;     boolean sortorder = true;     boolean category = true;     boolean field = true;     boolean filter = true;     string pc = null;     int productcount = 0;     string timestamp = null;      jsonobject obj = null;     try {         obj = new jsonobject(feed);     } catch (jsonexception e1) {         // todo auto-generated catch block         //e1.printstacktrace();      }      try {            searchterm = obj.getjsonobject("body").getstring("correctedword");             pageno = obj.getjsonobject("body").getstring("pageno");            sortorder = obj.getjsonobject("body").isnull("sortorder");            category = obj.getjsonobject("body").isnull("category");            field = obj.getjsonobject("body").isnull("field");            filter = obj.getjsonobject("body").getjsonobject("filter").isnull("filters");            pc = obj.getjsonobject("body").getstring("productcount").replaceall("[^\\d]", "");            productcount = integer.parseint(pc);            timestamp = (obj.getjsonobject("envelope").get("timestamp")).tostring().replaceall("[^\\d]", "");     } catch (jsonexception e) {         // todo auto-generated catch block         e.printstacktrace();      }      searchterm = searchterm.trim();      //condition eliminate pagination      if(!searchterm.isempty()){          if ((pageno.equals("1")) && (sortorder == true) && (category == true) && (field == true) && (filter == true)){              searchterm = searchterm.tolowercase();              system.out.println("in normalizer term : "+searchterm+","+timestamp+","+productcount);             system.out.println("entire json : "+feed);               collector.emit("stream1", new values(searchterm , timestamp , productcount ));              }      }       }    /**  * bolt emit field "word"   */ public void declareoutputfields(outputfieldsdeclarer declarer) {      declarer.declarestream("stream1", new fields("searchterm" ,"timestamp" ,"productcount"));  } } 

cassandrawriter bolt:

public class wordcounter extends basebasicbolt { static final logger logger = logger.getlogger(wordcounter.class); integer id; string name; map<string, integer> counters; cluster cluster ; session session ;  /**  * @ end of spout (when cluster shutdown  * show word counters  */ @override public void cleanup() {  }  public static session getsessionwithretry(cluster cluster, string keyspace) {         while (true) {             try {                 return cluster.connect(keyspace);             } catch (nohostavailableexception e) {                  utils.sleep(1000);             }         }      } public static cluster setupcassandraclient() {     return cluster.builder().addcontactpoint("192.168.1.229").build(); } /**  * on create   */ @override public void prepare(map stormconf, topologycontext context) {     this.counters = new hashmap<string, integer>();     this.name = context.getthiscomponentid();     this.id = context.getthistaskid();     cluster = setupcassandraclient();     session = wordcounter.getsessionwithretry(cluster,"query");      string query = "create table if not exists productcount(uid uuid primary key, "             + "term text , "             + "productcount varint,"              +"timestamp text );";      session.executeasync(query); }  @override public void declareoutputfields(outputfieldsdeclarer declarer) {}   @override public void execute(tuple input, basicoutputcollector collector) {     string term = input.getstring(0);     string timestamp = input.getstring(1);     int productcount = input.getinteger(2);      system.out.println("in counter : " +term+","+productcount+","+timestamp);     /**      * if word dosn't exist in map create      * this, if not add 1       */    string insertintotable = "insert productcount (uid, term, productcount, timestamp)"       + " values("+uuid.randomuuid()+","+"\'"+term+"\'"+","+productcount+","+"\'"+timestamp+"\'"+");" ;     session.executeasync(insertintotable);    } } 

please suggests me changes need do.

thanks in advnace !!

the error points missing ' in 1 of queries. problem if concatenate strings. avoid sort of problems should use prepared statements or @ least session.execute(query, params) (javadoc)


Comments

Popular posts from this blog

javascript - Google App Script ContentService downloadAsFile not working -

javascript - IIFE: var vs this - is there any difference? -

r - Grow a ffdf data frame on disk gradually -