Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,9 @@ public class DataTree {

private static final Logger LOG = LoggerFactory.getLogger(DataTree.class);

private static final int UTF8_SINGLE_BYTE_LIMIT = 0x80;
private static final int UTF8_TWO_BYTE_LIMIT = 0x800;

private final RateLogger RATE_LOGGER = new RateLogger(LOG, 15 * 60 * 1000);

/**
Expand Down Expand Up @@ -151,7 +154,30 @@ public class DataTree {
/**
* This hashtable lists the paths of the ephemeral nodes of a session.
*/
private final Map<Long, HashSet<String>> ephemerals = new ConcurrentHashMap<>();
private final Map<Long, EphemeralNodes> ephemerals = new ConcurrentHashMap<>();

private static final class EphemeralNodes {

private final HashSet<String> paths = new HashSet<>();
private volatile long pathBytes;

private void add(String path) {
synchronized (paths) {
if (paths.add(path)) {
pathBytes += getPathByteSize(path);
}
}
}

private void remove(String path) {
synchronized (paths) {
if (paths.remove(path)) {
pathBytes -= getPathByteSize(path);
}
}
}

}

/**
* This set contains the paths of all container nodes
Expand Down Expand Up @@ -189,14 +215,13 @@ public class DataTree {

private final DigestCalculator digestCalculator;

@SuppressWarnings("unchecked")
public Set<String> getEphemerals(long sessionId) {
HashSet<String> ret = ephemerals.get(sessionId);
EphemeralNodes ret = ephemerals.get(sessionId);
if (ret == null) {
return new HashSet<>();
}
synchronized (ret) {
return (HashSet<String>) ret.clone();
synchronized (ret.paths) {
return new HashSet<>(ret.paths);
}
}

Expand Down Expand Up @@ -225,11 +250,26 @@ public int getWatchCount() {
}

public int getEphemeralsCount() {
int result = 0;
for (HashSet<String> set : ephemerals.values()) {
result += set.size();
return ephemerals.values().stream().mapToInt(session -> session.paths.size()).sum();
}

/**
* Returns the maximum total size of a session's ephemeral paths in UTF-8 bytes,
* or zero if there are no ephemeral nodes. The total excludes
* {@code closeSession} serialization overhead.
*/
public long getMaxSessionEphemeralPathBytes() {
return ephemerals.values().stream().mapToLong(session -> session.pathBytes).max().orElse(0L);
}

private static long getPathByteSize(String path) {
long size = 0;
// Valid znode paths exclude surrogate characters, so each character uses at most three UTF-8 bytes.
for (int i = 0; i < path.length(); i++) {
char c = path.charAt(i);
size += c < UTF8_SINGLE_BYTE_LIMIT ? 1 : c < UTF8_TWO_BYTE_LIMIT ? 2 : 3;
}
return result;
return size;
}

/**
Expand Down Expand Up @@ -488,10 +528,7 @@ public void createNode(final String path, byte[] data, List<ACL> acl, long ephem
ttls.add(path);
ServerMetrics.getMetrics().TTL_NODE_CREATED_COUNT.add(1);
} else if (ephemeralOwner != 0) {
HashSet<String> list = ephemerals.computeIfAbsent(ephemeralOwner, k -> new HashSet<>());
synchronized (list) {
list.add(path);
}
ephemerals.computeIfAbsent(ephemeralOwner, k -> new EphemeralNodes()).add(path);
}
if (outputStat != null) {
child.copyStat(outputStat);
Expand Down Expand Up @@ -592,11 +629,9 @@ public void deleteNode(String path, long zxid) throws NoNodeException {
ttls.remove(path);
ServerMetrics.getMetrics().TTL_NODE_DELETED_COUNT.add(1);
} else if (owner != 0) {
Set<String> nodes = ephemerals.get(owner);
if (nodes != null) {
synchronized (nodes) {
nodes.remove(path);
}
EphemeralNodes session = ephemerals.get(owner);
if (session != null) {
session.remove(path);
}
}
}
Expand Down Expand Up @@ -969,7 +1004,7 @@ public ProcessTxnResult processTxn(TxnHeader header, Record txn, boolean isSubTx
long sessionId = header.getClientId();
if (txn != null) {
killSession(sessionId, header.getZxid(),
ephemerals.remove(sessionId),
removeEphemerals(sessionId),
((CloseSessionTxn) txn).getPaths2Delete());
} else {
killSession(sessionId, header.getZxid());
Expand Down Expand Up @@ -1144,7 +1179,12 @@ void killSession(long session, long zxid) {
// so there is no need for synchronization. The list is not
// changed here. Only create and delete change the list which
// are again called from FinalRequestProcessor in sequence.
killSession(session, zxid, ephemerals.remove(session), null);
killSession(session, zxid, removeEphemerals(session), null);
}

private Set<String> removeEphemerals(long sessionId) {
EphemeralNodes session = ephemerals.remove(sessionId);
return session == null ? null : session.paths;
}

void killSession(long session, long zxid, Set<String> paths2DeleteLocal,
Expand Down Expand Up @@ -1372,6 +1412,7 @@ public void deserialize(InputArchive ia, String tag) throws IOException {
aclCache.deserialize(ia);
nodes.clear();
pTrie.clear();
ephemerals.clear();
nodeDataSize.set(0);
String path = ia.readString("path");
while (!"/".equals(path)) {
Expand Down Expand Up @@ -1399,8 +1440,7 @@ public void deserialize(InputArchive ia, String tag) throws IOException {
} else if (ephemeralType == EphemeralType.TTL) {
ttls.add(path);
} else if (owner != 0) {
HashSet<String> list = ephemerals.computeIfAbsent(owner, k -> new HashSet<>());
list.add(path);
ephemerals.computeIfAbsent(owner, k -> new EphemeralNodes()).add(path);
}
}
path = ia.readString("path");
Expand Down Expand Up @@ -1472,10 +1512,10 @@ public synchronized WatchesSummary getWatchesSummary() {
*/
public void dumpEphemerals(PrintWriter writer) {
writer.println("Sessions with Ephemerals (" + ephemerals.keySet().size() + "):");
for (Entry<Long, HashSet<String>> entry : ephemerals.entrySet()) {
for (Entry<Long, EphemeralNodes> entry : ephemerals.entrySet()) {
writer.print("0x" + Long.toHexString(entry.getKey()));
writer.println(":");
Set<String> tmp = entry.getValue();
Set<String> tmp = entry.getValue().paths;
if (tmp != null) {
synchronized (tmp) {
for (String path : tmp) {
Expand All @@ -1498,9 +1538,9 @@ public void shutdownWatcher() {
*/
public Map<Long, Set<String>> getEphemerals() {
Map<Long, Set<String>> ephemeralsCopy = new HashMap<>();
for (Entry<Long, HashSet<String>> e : ephemerals.entrySet()) {
synchronized (e.getValue()) {
ephemeralsCopy.put(e.getKey(), new HashSet<>(e.getValue()));
for (Entry<Long, EphemeralNodes> e : ephemerals.entrySet()) {
synchronized (e.getValue().paths) {
ephemeralsCopy.put(e.getKey(), new HashSet<>(e.getValue().paths));
}
}
return ephemeralsCopy;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,7 @@ public class ZooKeeperServer implements SessionExpirer, ServerStats.Provider {
// Add a enable/disable option for now, we should remove this one when
// this feature is confirmed to be stable
public static final String CLOSE_SESSION_TXN_ENABLED = "zookeeper.closeSessionTxn.enabled";
private static final String MAX_SESSION_EPHEMERAL_PATH_BYTES = "max_session_ephemeral_path_bytes";
private static boolean closeSessionTxnEnabled = true;
private volatile CountDownLatch restoreLatch;
// exclusive lock for taking snapshot and restore
Expand Down Expand Up @@ -2023,6 +2024,8 @@ protected void registerMetrics() {

rootContext.registerGauge("watch_count", zkdb.getDataTree()::getWatchCount);
rootContext.registerGauge("ephemerals_count", zkdb.getDataTree()::getEphemeralsCount);
rootContext.registerGauge(MAX_SESSION_EPHEMERAL_PATH_BYTES,
() -> getZKDatabase().getDataTree().getMaxSessionEphemeralPathBytes());

rootContext.registerGauge("approximate_data_size", zkdb.getDataTree()::cachedApproximateDataSize);

Expand Down Expand Up @@ -2073,6 +2076,7 @@ protected void unregisterMetrics() {

rootContext.unregisterGauge("watch_count");
rootContext.unregisterGauge("ephemerals_count");
rootContext.unregisterGauge(MAX_SESSION_EPHEMERAL_PATH_BYTES);
rootContext.unregisterGauge("approximate_data_size");

rootContext.unregisterGauge("global_sessions");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -581,6 +581,7 @@ public CommandResponse runGet(ZooKeeperServer zkServer, Map<String, String> kwar
* - "znode_count": Integer
* - "watch_count": Integer
* - "ephemerals_count": Integer
* - "max_session_ephemeral_path_bytes": Number (UTF-8 bytes)
* - "approximate_data_size": Long
* - "open_file_descriptor_count": Long (unix only)
* - "max_file_descriptor_count": Long (unix only)
Expand Down
Loading