jgroupsmember.java
来自「JGRoups源码」· Java 代码 · 共 468 行 · 第 1/2 页
JAVA
468 行
/**
* A MembershipListener callback. It is called when a member is suspected being crashed,
* but has not yet been excluded from the group.
* @param suspectedMbr is the Address of the suspected crashed member.
*/
public void suspect(Address suspectedMbr) {
if(logger.isDebugEnabled())
logger.debug("[" + getLocalAddress() + "] MembershipListener.suspect(Address=" +
suspectedMbr + ") is called notifying a suspected crashed member...");
}
/**
* A MembershipListener callback. It is called when a change in group membership has
* occurred, either new member joins, existing member leaves or has been suspected crashed.
* @param mbrView is a View containing Address of all members in the group.
*/
public void viewAccepted(View mbrView) {
String addr=getLocalAddress();
System.out.println("[" + addr + "] MembershipListener.viewAccepted(View) is called...");
if(null == mbrView) {
logger.warn("[" + addr + "] a null View is received.");
return;
}
Vector mbrs=mbrView.getMembers();
System.out.println("[" + addr + "] View has " + mbrs.size() + " members.");
Address mbr=null;
for(Iterator it=mbrs.iterator(); it.hasNext();) {
mbr=(Address)it.next();
logger.debug("[" + addr + "] View a member address: " + mbr);
}
}
/**
* A MessageListener callback. This method is called when another member invokes its
* Channel.getState() requesting for group state, and this member is chosen to respond,
* which would indicate that this member is a group coordinator.
* @return It returns a byte[] containing a serialized Address of the server known
* by this member. It returns null when failed to serialize the Address.
*/
public byte[] getState() {
String addr=getLocalAddress();
logger.info("[" + addr + "] is responding to a Channel.getState() inquiry from another" +
" member...");
try {
byte[] result=null;
logger.debug("[" + addr + "] waiting to lock serverAddress to read...");
synchronized(serverAddressMutex) {
result=Util.objectToByteBuffer(serverAddress);
logger.debug("[" + addr + "] has read serverAddress;" +
" releasing serverAddress lock...");
}
logger.info("[" + addr + "] replying Channel.getState() inquiry...");
return result;
}
catch(Exception ex) {
logger.error("[" + addr + "] failed to serialize reply to the Channel.getState()" +
" inquiry; Exception: ", ex);
logger.error("[" + addr + "] replying null to Channel.getState() inquiry...");
return null;
}
}
/**
* A MessageListener callback. This method is called to receive a group state after this
* member has invoked a Channel.getState() call. The state received is a serialized
* Address of server known by the current group coordinator.
* @param state is a byte[] containing an Address of server known by the current group
* coordinator. It could be null if the group coordinator failed to
* serialize the server address, or the coordinator does not have the
* server address when responding to this member's inquiry.
*/
public void setState(byte[] state) {
String addr=getLocalAddress();
logger.info("[" + addr + "] is receiving a group state...");
try {
if(null != state) {
logger.debug("[" + addr + "] waiting to lock serverAddress to write...");
synchronized(serverAddressMutex) {
serverAddress=(Address)Util.objectFromByteBuffer(state);
logger.debug("[" + addr + "] has written to serverAddress;" +
" releasing serverAddress lock...");
}
if(logger.isDebugEnabled())
logger.debug("[" + addr + "] Got a state from a group coordinator that " +
"the server address is " + serverAddress);
}
else
logger.error("[" + addr + "] Got null state.");
}
catch(Exception ex) {
logger.error("[" + addr + "] failed to de-serialize the state; Exception: ", ex);
}
}
/**
* A MessageListener callback. It is called to receive a Message sent from another member.
* @param message the Message received.
*/
public void receive(Message message) {
String addr=getLocalAddress();
if(logger.isDebugEnabled())
logger.debug("[" + addr + "] is receiving a Message via MessageListener...");
if(null == message) {
logger.error("[" + addr + "] got a null Message.");
return;
}
Address src=message.getSrc();
if(message.getObject() instanceof String) {
String contents=(String)message.getObject();
if(logger.isDebugEnabled()) {
logger.debug("[" + addr + "] got a Message \"" + contents + "\" from member at " +
src.toString());
}
}
else {
if(logger.isDebugEnabled()) {
logger.debug("[" + addr + "] got a Message from member at " + src.toString());
}
}
}
/**
* A RequestHandler callback. It is called to receive a Message multi-casted by another
* member in the group, e.g. by a MessageDispatcher.castMessage(Message).
* @param message the Message received.
* @return If it is a server, it returns a String "ACK". If it is a client,
* it returns null.
*/
public Object handle(Message message) {
String addr=getLocalAddress();
if(logger.isDebugEnabled())
logger.debug("[" + addr + "] is receiving a Message via RequestHandler...");
if(null == message) {
logger.error("[" + addr + "] got a null Message.");
return null;
}
Address src=message.getSrc();
if(message.getObject() instanceof String) {
String contents=(String)message.getObject();
if(logger.isDebugEnabled()) {
logger.debug("[" + addr + "] got a Message \"" + contents + "\" from member at " +
src.toString());
}
if(contents.equals("server")) {
logger.debug("[" + addr + "] waiting to lock serverAddress to write...");
synchronized(serverAddressMutex) {
serverAddress=src;
if(logger.isDebugEnabled())
logger.debug("[" + addr + "] has written to serverAddress = " + src +
"; releasing serverAddress lock..");
}
}
else if(contents.equals("stop")) {
logger.debug("[" + addr + "] received a \"stop\" message;" +
" waiting to lock serverEnds...");
synchronized(serverEnds) {
logger.debug("[" + addr + "] waking up server main thread...");
serverEnds.notifyAll();
}
}
}
else {
if(logger.isDebugEnabled()) {
logger.debug("[" + addr + "] got a Message from member at " + src.toString());
}
}
if(role.equals("server")) {
logger.debug("[" + addr + "] replying \"ACK\"...");
return "ACK";
}
else {
logger.debug("[" + addr + "] replying null...");
return null;
}
}
/*
* Returns a String of the Address this member is connected to the JGroups. It returns
* null if the channel is not connected.
*/
private String getLocalAddress() {
if(null == channel)
return null;
else {
Address addr=channel.getLocalAddress();
if(null == addr) // the channel is in closed state
return null;
else
return addr.toString();
}
}
public static void main(String args[]) {
String role, propFile, message;
int numOfRepeats;
JGroupsMember member=null;
if(4 != args.length) {
System.out.println(
"Usage: JGroupsMember <role> <configXmlFile> <message> <repeats>");
System.out.println(" role = client or server.");
System.out.println(" configXmlFile = the JGroups protocol stack config file.");
System.out.println(" For example: config/total-token.xml");
System.out.println(" message = the message to multi-cast to the group.");
System.out.println(" A message \"stop\" from a client will stop");
System.out.println(" all servers.");
System.out.println(" repeats = if the role is \"client\", then how many");
System.out.println(" times to repeat the sending of the message.");
return;
}
else {
role=args[0];
propFile=args[1];
message=args[2];
numOfRepeats=Integer.parseInt(args[3]);
}
member=new JGroupsMember(role, propFile, message, numOfRepeats);
member.init();
member.run();
}
}
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?