@@ -84,7 +84,7 @@ public void setUp() throws IOException {
84
84
try {client .unsubscribe (HOST , PORT , "Trades1" , "subtrades2" );}catch (Exception e ){}
85
85
try {client .unsubscribe (HOST , PORT , "Trades1" );}catch (Exception e ){}
86
86
try {client .unsubscribe (HOST , PORT , "Trades" , "subTread1" );}catch (Exception e ){}
87
- clear_env ();
87
+ try { clear_env ();} catch ( Exception e ){}
88
88
conn .run ("st2 = streamTable(1000000:0,`tag`ts`data,[INT,TIMESTAMP,DOUBLE])\n " +
89
89
"enableTableShareAndPersistence(table=st2, tableName=`Trades1, asynWrite=true, compress=true, cacheSize=20000, retentionMinutes=180)\t \n " );
90
90
}
@@ -2207,23 +2207,23 @@ public void run() {
2207
2207
thread1 .start ();
2208
2208
Thread .sleep (2000 );
2209
2209
thread2 .start ();
2210
- // List<IMessage> messages = poller.poll(1000,1000);
2211
- // System.out.println("messages" + messages.size());
2212
- // MessageHandler_handler1(messages);
2213
- // Thread.sleep(1000);
2210
+ List <IMessage > messages = poller .poll (1000 ,1000 );
2211
+ System .out .println ("messages" + messages .size ());
2212
+ MessageHandler_handler1 (messages );
2213
+ Thread .sleep (1000 );
2214
2214
thread .join ();
2215
2215
Thread .sleep (10000 );
2216
2216
controller_conn .run ("try{startDataNode('" +HOST +":" +port_list [1 ]+"')}catch(ex){}" );
2217
2217
//Thread.sleep(1000);
2218
2218
List <IMessage > messages1 = poller .poll (1000 ,1000 );
2219
2219
Thread .sleep (1000 );
2220
2220
System .out .println (messages1 .size ());
2221
- Assert .assertEquals (1000 ,messages1 .size ());
2221
+ // Assert.assertEquals(1000,messages1.size());
2222
2222
MessageHandler_handler1 (messages1 );
2223
2223
Thread .sleep (1000 );
2224
2224
BasicTable re = (BasicTable )conn .run ("select tag ,now,deltas(now) from Receive order by deltas(now) desc \n " );
2225
2225
System .out .println (re .getString ());
2226
- // Assert.assertEquals(1000,re.rows());
2226
+ Assert .assertEquals (1000 ,re .rows ());
2227
2227
Assert .assertEquals (true ,Integer .valueOf (re .getColumn (2 ).get (0 ).toString ())>1000 );
2228
2228
DBConnection conn2 = new DBConnection ();
2229
2229
conn2 .connect (HOST ,port_list [1 ],"admin" ,"123456" );
@@ -2281,7 +2281,7 @@ public void test_PollingClient_subscribe_resubTimeout_subOnce_not_set() throws I
2281
2281
Thread .sleep (20000 );
2282
2282
conn3 .run ("n=3000;t=table(1..n as tag,timestamp(1..n) as ts,take(100.0,n) as data);" + "Trades.append!(t)" );
2283
2283
Thread .sleep (5000 );
2284
- List <IMessage > messages2 = poller .poll (3000 ,1000 );
2284
+ List <IMessage > messages2 = poller .poll (3000 ,3000 );
2285
2285
MessageHandler_handler (messages2 );
2286
2286
controller_conn .run ("try{startDataNode('" +HOST +":" +port_list [2 ]+"')}catch(ex){}" );
2287
2287
Thread .sleep (5000 );
0 commit comments