virajjasani commented on code in PR #8584:
URL: https://github.com/apache/hbase/pull/8584#discussion_r4180818258
##########
hbase-server/src/main/java/org/apache/hadoop/hbase/master/ServerManager.java:
##########
@@ -289,35 +289,50 @@ ServerName regionServerStartup(RegionServerStartupRequest
request, int versionNu
private void updateLastFlushedSequenceIds(ServerName sn, ServerMetrics hsl) {
for (Entry<byte[], RegionMetrics> entry :
hsl.getRegionMetrics().entrySet()) {
byte[] encodedRegionName =
Bytes.toBytes(RegionInfo.encodeRegionName(entry.getKey()));
- Long existingValue = flushedSequenceIdByRegion.get(encodedRegionName);
- long l = entry.getValue().getCompletedSequenceId();
- // Don't let smaller sequence ids override greater sequence ids.
- if (LOG.isTraceEnabled()) {
- LOG.trace(Bytes.toString(encodedRegionName) + ", existingValue=" +
existingValue
- + ", completeSequenceId=" + l);
- }
- if (existingValue == null || (l != HConstants.NO_SEQNUM && l >
existingValue)) {
- flushedSequenceIdByRegion.put(encodedRegionName, l);
- } else if (l != HConstants.NO_SEQNUM && l < existingValue) {
- LOG.warn("RegionServer " + sn + " indicates a last flushed sequence id
(" + l
- + ") that is less than the previous last flushed sequence id (" +
existingValue
- + ") for region " + Bytes.toString(entry.getKey()) + " Ignoring.");
- }
+ final long completedSeqId = entry.getValue().getCompletedSequenceId();
+ // Atomic read-modify-write so a concurrent reportRegionOpen seed (which
uses
+ // merge(Math::max)) cannot be clobbered by a stale in-flight heartbeat
carrying a
+ // lower completedSequenceId. Don't let smaller sequence ids override
greater ones.
+ flushedSequenceIdByRegion.compute(encodedRegionName, (k, existingValue)
-> {
+ if (LOG.isTraceEnabled()) {
+ LOG.trace(Bytes.toString(k) + ", existingValue=" + existingValue +
", completeSequenceId="
+ + completedSeqId);
+ }
+ if (existingValue == null) {
+ return completedSeqId;
+ }
+ if (completedSeqId != HConstants.NO_SEQNUM && completedSeqId >
existingValue) {
+ return completedSeqId;
+ }
+ if (completedSeqId != HConstants.NO_SEQNUM && completedSeqId <
existingValue) {
+ LOG.warn("RegionServer " + sn + " indicates a last flushed sequence
id (" + completedSeqId
+ + ") that is less than the previous last flushed sequence id (" +
existingValue
+ + ") for region " + Bytes.toString(entry.getKey()) + " Ignoring.");
+ }
+ return existingValue;
+ });
ConcurrentNavigableMap<byte[], Long> storeFlushedSequenceId =
computeIfAbsent(storeFlushedSequenceIdsByRegion, encodedRegionName,
() -> new ConcurrentSkipListMap<>(Bytes.BYTES_COMPARATOR));
for (Entry<byte[], Long> storeSeqId :
entry.getValue().getStoreSequenceId().entrySet()) {
byte[] family = storeSeqId.getKey();
- existingValue = storeFlushedSequenceId.get(family);
- l = storeSeqId.getValue();
- if (LOG.isTraceEnabled()) {
- LOG.trace(Bytes.toString(encodedRegionName) + ", family=" +
Bytes.toString(family)
- + ", existingValue=" + existingValue + ", completeSequenceId=" +
l);
- }
- // Don't let smaller sequence ids override greater sequence ids.
- if (existingValue == null || (l != HConstants.NO_SEQNUM && l >
existingValue.longValue())) {
- storeFlushedSequenceId.put(family, l);
- }
+ final long storeCompletedSeqId = storeSeqId.getValue();
+ storeFlushedSequenceId.compute(family, (k, existingValue) -> {
+ if (LOG.isTraceEnabled()) {
+ LOG.trace(Bytes.toString(encodedRegionName) + ", family=" +
Bytes.toString(k)
+ + ", existingValue=" + existingValue + ", completeSequenceId=" +
storeCompletedSeqId);
+ }
+ if (existingValue == null) {
+ return storeCompletedSeqId;
+ }
+ if (
+ storeCompletedSeqId != HConstants.NO_SEQNUM
+ && storeCompletedSeqId > existingValue.longValue()
Review Comment:
`storeCompletedSeqId > existingValue` is also good enough
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]