diff --git a/src/java/org/apache/cassandra/net/http/HttpRequestHandler.java b/src/java/org/apache/cassandra/net/http/HttpRequestHandler.java index f7e3221518..f174e36409 100644 --- a/src/java/org/apache/cassandra/net/http/HttpRequestHandler.java +++ b/src/java/org/apache/cassandra/net/http/HttpRequestHandler.java @@ -312,7 +312,7 @@ public class HttpRequestHandler implements Runnable EndPoint[] liveNodes = liveNodeList.toArray(new EndPoint[0]); Arrays.sort(liveNodes); - String[] sHeaders = {"Node No.", "Host:Port", "Status", "Leader", "Load Info", "Token", "Generation No."}; + String[] sHeaders = {"Node No.", "Host:Port", "Status", "Load Info", "Token", "Generation No."}; formatter.startTable(); formatter.addHeaders(sHeaders); int iNodeNumber = 0; @@ -328,9 +328,6 @@ public class HttpRequestHandler implements Runnable //Status String status = ( FailureDetector.instance().isAlive(curNode) ) ? "Up" : "Down"; formatter.addCol(status); - //Leader - boolean isLeader = StorageService.instance().isLeader(curNode); - formatter.addCol(Boolean.toString(isLeader)); //Load Info String loadInfo = getLoadInfo(curNode); formatter.addCol(loadInfo); diff --git a/src/java/org/apache/cassandra/service/LeaderElector.java b/src/java/org/apache/cassandra/service/LeaderElector.java deleted file mode 100644 index 3041233fc4..0000000000 --- a/src/java/org/apache/cassandra/service/LeaderElector.java +++ /dev/null @@ -1,266 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.service; - -import java.io.IOException; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.SortedMap; -import java.util.TreeMap; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; -import java.util.concurrent.locks.Condition; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; - -import org.apache.cassandra.concurrent.DebuggableThreadPoolExecutor; -import org.apache.cassandra.concurrent.ThreadFactoryImpl; -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.gms.ApplicationState; -import org.apache.cassandra.gms.EndPointState; -import org.apache.cassandra.gms.Gossiper; -import org.apache.cassandra.gms.IEndPointStateChangeSubscriber; -import org.apache.cassandra.net.EndPoint; -import org.apache.cassandra.utils.LogUtil; -import org.apache.log4j.Logger; -import org.apache.zookeeper.*; -import org.apache.zookeeper.ZooDefs.Ids; -import org.apache.zookeeper.data.Stat; - -class LeaderElector implements IEndPointStateChangeSubscriber -{ - private static Logger logger_ = Logger.getLogger(LeaderElector.class); - protected static final String leaderState_ = "LEADER"; - private static LeaderElector instance_ = null; - private static Lock createLock_ = new ReentrantLock(); - - /* - * Factory method that gets an instance of the StorageService - * class. - */ - public static LeaderElector instance() - { - if ( instance_ == null ) - { - LeaderElector.createLock_.lock(); - try - { - if ( instance_ == null ) - { - instance_ = new LeaderElector(); - } - } - finally - { - createLock_.unlock(); - } - } - return instance_; - } - - /* The elected leader. */ - private AtomicReference leader_; - private Condition condition_; - private ExecutorService leaderElectionService_ = new DebuggableThreadPoolExecutor("LEADER-ELECTOR"); - - private class LeaderDeathMonitor implements Runnable - { - private String pathCreated_; - - LeaderDeathMonitor(String pathCreated) - { - pathCreated_ = pathCreated; - } - - public void run() - { - ZooKeeper zk = StorageService.instance().getZooKeeperHandle(); - String path = "/Cassandra/" + DatabaseDescriptor.getClusterName() + "/Leader"; - try - { - String createPath = path + "/L-"; - LeaderElector.createLock_.lock(); - while( true ) - { - /* Get all znodes under the Leader znode */ - List values = zk.getChildren(path, false); - SortedMap suffixToZnode = getSuffixToZnodeMapping(values); - String value = suffixToZnode.get( suffixToZnode.firstKey() ); - /* - * Get the first znode and if it is the - * pathCreated created above then the data - * in that znode is the leader's identity. - */ - if ( leader_ == null ) - { - leader_ = new AtomicReference( EndPoint.fromBytes( zk.getData(path + "/" + value, false, null) ) ); - } - else - { - leader_.set( EndPoint.fromBytes( zk.getData(path + "/" + value, false, null) ) ); - /* Disseminate the state as to who the leader is. */ - onLeaderElection(); - } - logger_.debug("Elected leader is " + leader_ + " @ znode " + ( path + "/" + value ) ); - /* We need only the last portion of this znode */ - int index = getLocalSuffix(); - if ( index > suffixToZnode.firstKey() ) - { - String pathToCheck = path + "/" + getImmediatelyPrecedingZnode(suffixToZnode, index); - Stat stat = zk.exists(pathToCheck, true); - if ( stat != null ) - { - logger_.debug("Awaiting my turn ..."); - condition_.await(); - logger_.debug("Checking to see if leader is around ..."); - } - } - else - { - break; - } - } - } - catch ( InterruptedException ex ) - { - logger_.warn(LogUtil.throwableToString(ex)); - } - catch ( IOException ex ) - { - logger_.warn(LogUtil.throwableToString(ex)); - } - catch ( KeeperException ex ) - { - logger_.warn(LogUtil.throwableToString(ex)); - } - finally - { - LeaderElector.createLock_.unlock(); - } - } - - private SortedMap getSuffixToZnodeMapping(List values) - { - SortedMap suffixZNodeMap = new TreeMap(); - for ( String znode : values ) - { - String[] peices = znode.split("-"); - suffixZNodeMap.put(Integer.parseInt( peices[1] ), znode); - } - return suffixZNodeMap; - } - - private String getImmediatelyPrecedingZnode(SortedMap suffixToZnode, int index) - { - List suffixes = new ArrayList( suffixToZnode.keySet() ); - Collections.sort(suffixes); - int position = Collections.binarySearch(suffixes, index); - return suffixToZnode.get( suffixes.get( position - 1 ) ); - } - - /** - * If the local node's leader related znode is L-3 - * this method will return 3. - * @return suffix portion of L-3. - */ - private int getLocalSuffix() - { - String[] peices = pathCreated_.split("/"); - String leaderPeice = peices[peices.length - 1]; - String[] leaderPeices = leaderPeice.split("-"); - return Integer.parseInt( leaderPeices[1] ); - } - } - - private LeaderElector() - { - condition_ = LeaderElector.createLock_.newCondition(); - } - - /** - * Use to inform interested parties about the change in the state - * for specified endpoint - * - * @param endpoint endpoint for which the state change occured. - * @param epState state that actually changed for the above endpoint. - */ - public void onChange(EndPoint endpoint, EndPointState epState) - { - /* node identifier for this endpoint on the identifier space */ - ApplicationState leaderState = epState.getApplicationState(LeaderElector.leaderState_); - if (leaderState != null && !leader_.equals(endpoint)) - { - EndPoint leader = EndPoint.fromString( leaderState.getState() ); - logger_.debug("New leader in the cluster is " + leader); - leader_.set(endpoint); - } - } - - void start() throws Throwable - { - /* Register with the Gossiper for EndPointState notifications */ - Gossiper.instance().register(this); - logger_.debug("Starting the leader election process ..."); - ZooKeeper zk = StorageService.instance().getZooKeeperHandle(); - String path = "/Cassandra/" + DatabaseDescriptor.getClusterName() + "/Leader"; - String createPath = path + "/L-"; - - /* Create the znodes under the Leader znode */ - logger_.debug("Attempting to create znode " + createPath); - String pathCreated = zk.create(createPath, EndPoint.toBytes( StorageService.getLocalControlEndPoint() ), Ids.OPEN_ACL_UNSAFE, (CreateMode.EPHEMERAL_SEQUENTIAL) ); - logger_.debug("Created znode under leader znode " + pathCreated); - leaderElectionService_.submit(new LeaderDeathMonitor(pathCreated)); - } - - void signal() - { - logger_.debug("Signalling others to check on leader ..."); - try - { - LeaderElector.createLock_.lock(); - condition_.signal(); - } - finally - { - LeaderElector.createLock_.unlock(); - } - } - - EndPoint getLeader() - { - return (leader_ != null ) ? leader_.get() : StorageService.getLocalStorageEndPoint(); - } - - private void onLeaderElection() throws InterruptedException, IOException - { - /* - * If the local node is the leader then not only does he - * diseminate the information but also starts the M/R job - * tracker. Non leader nodes start the M/R task tracker - * thereby initializing the M/R subsystem. - */ - if ( StorageService.instance().isLeader(leader_.get()) ) - { - Gossiper.instance().addApplicationState(LeaderElector.leaderState_, new ApplicationState(leader_.toString())); - } - } -} diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 8840685441..aa3d2295e9 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -317,91 +317,12 @@ public final class StorageService implements IEndPointStateChangeSubscriber, Sto else nodePicker_ = new RackUnawareStrategy(tokenMetadata_, partitioner_, DatabaseDescriptor.getReplicationFactor(), DatabaseDescriptor.getStoragePort()); } - - private void reportToZookeeper() throws Throwable - { - try - { - zk_ = new ZooKeeper(DatabaseDescriptor.getZkAddress(), DatabaseDescriptor.getZkSessionTimeout(), new Watcher() - { - public void process(WatchedEvent we) - { - String path = "/Cassandra/" + DatabaseDescriptor.getClusterName() + "/Leader"; - String eventPath = we.getPath(); - logger_.debug("PROCESS EVENT : " + eventPath); - if (eventPath != null && (eventPath.contains(path))) - { - logger_.debug("Signalling the leader instance ..."); - LeaderElector.instance().signal(); - } - } - }); - - Stat stat = zk_.exists("/", false); - if ( stat != null ) - { - stat = zk_.exists("/Cassandra", false); - if ( stat == null ) - { - logger_.debug("Creating the Cassandra znode ..."); - zk_.create("/Cassandra", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - - String path = "/Cassandra/" + DatabaseDescriptor.getClusterName(); - stat = zk_.exists(path, false); - if ( stat == null ) - { - logger_.debug("Creating the cluster znode " + path); - zk_.create(path, new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - - /* Create the Leader, Locks and Misc znode */ - stat = zk_.exists(path + "/Leader", false); - if ( stat == null ) - { - logger_.debug("Creating the leader znode " + path); - zk_.create(path + "/Leader", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - - stat = zk_.exists(path + "/Locks", false); - if ( stat == null ) - { - logger_.debug("Creating the locks znode " + path); - zk_.create(path + "/Locks", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - - stat = zk_.exists(path + "/Misc", false); - if ( stat == null ) - { - logger_.debug("Creating the misc znode " + path); - zk_.create(path + "/Misc", new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - } - } - } - catch ( KeeperException ke ) - { - LogUtil.throwableToString(ke); - /* do the re-initialize again. */ - reportToZookeeper(); - } - } - + protected ZooKeeper getZooKeeperHandle() { return zk_; } - public boolean isLeader(EndPoint endpoint) - { - EndPoint leader = getLeader(); - return leader.equals(endpoint); - } - - public EndPoint getLeader() - { - return LeaderElector.instance().getLeader(); - } - public void registerComponentForShutdown(IComponentShutdown component) { components_.add(component); @@ -440,14 +361,6 @@ public final class StorageService implements IEndPointStateChangeSubscriber, Sto /* starts a load timer thread */ loadTimer_.schedule( new LoadDisseminator(), StorageService.threshold_, StorageService.threshold_); - /* report our existence to ZooKeeper instance and start the leader election service */ - - //reportToZookeeper(); - /* start the leader election algorithm */ - //LeaderElector.instance().start(); - /* start the map reduce framework */ - //startMapReduceFramework(); - /* Start the storage load balancer */ storageLoadBalancer_.start(); /* Register with the Gossiper for EndPointState notifications */ diff --git a/src/java/org/apache/cassandra/service/ZookeeperWatcher.java b/src/java/org/apache/cassandra/service/ZookeeperWatcher.java deleted file mode 100644 index 79e9ac85a6..0000000000 --- a/src/java/org/apache/cassandra/service/ZookeeperWatcher.java +++ /dev/null @@ -1,44 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.cassandra.service; - -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.log4j.Logger; -import org.apache.zookeeper.WatchedEvent; -import org.apache.zookeeper.Watcher; - - -public class ZookeeperWatcher implements Watcher -{ - private static final Logger logger_ = Logger.getLogger(ZookeeperWatcher.class); - private static final String leader_ = "/Cassandra/" + DatabaseDescriptor.getClusterName() + "/Leader"; - private static final String lock_ = "/Cassandra/" + DatabaseDescriptor.getClusterName() + "/Locks"; - - public void process(WatchedEvent we) - { - String eventPath = we.getPath(); - logger_.debug("PROCESS EVENT : " + eventPath); - if (eventPath != null && (eventPath.contains(leader_))) - { - logger_.debug("Signalling the leader instance ..."); - LeaderElector.instance().signal(); - } - - } -}