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
Post a Comment