Skip to content

Commit a122d14

Browse files
author
chengyitian
committed
AJ-692: fix issue about users set when subscribe;
1 parent 9b7dbcc commit a122d14

File tree

1 file changed

+2
-2
lines changed

1 file changed

+2
-2
lines changed

src/com/xxdb/streaming/client/AbstractClient.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -591,7 +591,6 @@ protected BlockingQueue<List<IMessage>> subscribeInternal(String host, int port,
591591

592592
List<String> tp = Arrays.asList(site.host, String.valueOf(site.port), tableName, actionName);
593593
List<String> usr = Arrays.asList(userName, passWord);
594-
users.put(tp, usr);
595594

596595
params.add(new BasicString(localIP));
597596
params.add(new BasicInt(this.listeningPort));
@@ -609,6 +608,7 @@ protected BlockingQueue<List<IMessage>> subscribeInternal(String host, int port,
609608

610609
re = dbConn.run("publishTable", params);
611610
connList.add(dbConn);
611+
users.put(tp, usr);
612612
} catch (IOException e) {
613613
log.error("Connect to site " + site.host + ":" + site.port + " failed: " + e.getMessage());
614614
}
@@ -655,7 +655,6 @@ protected BlockingQueue<List<IMessage>> subscribeInternal(String host, int port,
655655
checkServerVersion(host, port);
656656
List<String> tp = Arrays.asList(host, String.valueOf(port), tableName, actionName);
657657
List<String> usr = Arrays.asList(userName, passWord);
658-
users.put(tp, usr);
659658

660659
dbConn = createSubscribeInternalDBConnection();
661660
subscribeInternalConnect(dbConn, host, port, userName, passWord);
@@ -700,6 +699,7 @@ protected BlockingQueue<List<IMessage>> subscribeInternal(String host, int port,
700699

701700
re = dbConn.run("publishTable", params);
702701
connList.add(dbConn);
702+
users.put(tp, usr);
703703
if (ifUseBackupSite) {
704704
synchronized (subInfos_){
705705
subInfos_.put(topic, deserializer);

0 commit comments

Comments
 (0)