From 0251a8f0f66055960c2dc5502b6496313bf074f3 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 28 Apr 2011 13:36:07 +0000 Subject: [PATCH] trigger read repair correctly forLOCAL_QUORUM reads patch by jbellis; reviewed by slebresne for CASSANDRA-2556 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1097455 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../service/DatacenterReadCallback.java | 30 +++++-------------- .../cassandra/service/ReadCallback.java | 27 +++++++++++++++-- 3 files changed, 34 insertions(+), 24 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 8e79dcd720..ac99d927f0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -3,6 +3,7 @@ * move gossip heartbeat back to its own thread (CASSANDRA-2554) * fix incorrect use of NBHM.size in ReadCallback that could cause reads to time out even when responses were received (CASSAMDRA-2552) + * trigger read repair correctly for LOCAL_QUORUM reads (CASSANDRA-2556) 0.7.5 diff --git a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java index 53323d3692..3e6a37ba99 100644 --- a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java +++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java @@ -48,31 +48,17 @@ public class DatacenterReadCallback extends ReadCallback } @Override - public void response(Message message) + protected boolean waitingFor(Message message) { - resolver.preprocess(message); - - int n = localdc.equals(snitch.getDatacenter(message.getFrom())) - ? received.incrementAndGet() - : received.get(); - - if (n == blockfor && resolver.isDataPresent()) - { - condition.signal(); - maybeResolveForRepair(); - } + return localdc.equals(snitch.getDatacenter(message.getFrom())); } - - @Override - public void response(ReadResponse result) - { - ((RowDigestResolver) resolver).injectPreProcessed(result); - if (received.incrementAndGet() == blockfor && resolver.isDataPresent()) - { - condition.signal(); - maybeResolveForRepair(); - } + @Override + protected boolean waitingFor(ReadResponse response) + { + // cheat and leverage our knowledge that a local read is the only way the ReadResponse + // version of this method gets called + return true; } @Override diff --git a/src/java/org/apache/cassandra/service/ReadCallback.java b/src/java/org/apache/cassandra/service/ReadCallback.java index c482084ab0..a3e929c01b 100644 --- a/src/java/org/apache/cassandra/service/ReadCallback.java +++ b/src/java/org/apache/cassandra/service/ReadCallback.java @@ -128,17 +128,40 @@ public class ReadCallback implements IAsyncCallback public void response(Message message) { resolver.preprocess(message); - if (received.incrementAndGet() >= blockfor && resolver.isDataPresent()) + int n = waitingFor(message) + ? received.incrementAndGet() + : received.get(); + if (n >= blockfor && resolver.isDataPresent()) { condition.signal(); maybeResolveForRepair(); } } + /** + * @return true if the message counts towards the blockfor threshold + * TODO turn the Message into a response so we don't need two versions of this method + */ + protected boolean waitingFor(Message message) + { + return true; + } + + /** + * @return true if the response counts towards the blockfor threshold + */ + protected boolean waitingFor(ReadResponse response) + { + return true; + } + public void response(ReadResponse result) { ((RowDigestResolver) resolver).injectPreProcessed(result); - if (received.incrementAndGet() >= blockfor && resolver.isDataPresent()) + int n = waitingFor(result) + ? received.incrementAndGet() + : received.get(); + if (n >= blockfor && resolver.isDataPresent()) { condition.signal(); maybeResolveForRepair();