From b8efe711e406a367c3632b608b80ab09ac965310 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 12 Sep 2022 15:09:08 -0700 Subject: [PATCH 1/9] KAFKA-14214: Introduce read-write lock to StandardAuthorizer for consistent ACL reads. The issue with StandardAuthorizer#authorize is, that it looks up aclsByResources (which is of type ConcurrentSkipListMap)twice for every authorize call and uses Iterator with weak consistency guarantees on top of aclsByResources. This can cause the authorize function call to process the concurrent writes out of order. Implemented ReadWrite lock at StandardAuthorizer level to make sure the reads are strongly consistent with write order. --- .../authorizer/StandardAuthorizer.java | 74 +++++++++++++------ .../authorizer/StandardAuthorizerData.java | 66 +++++------------ 2 files changed, 69 insertions(+), 71 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java index 42f03367c22e8..d8492590d5e58 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java @@ -39,6 +39,8 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; +import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.Supplier; import static org.apache.kafka.server.authorizer.AuthorizationResult.ALLOWED; import static org.apache.kafka.server.authorizer.AuthorizationResult.DENIED; @@ -58,19 +60,39 @@ public class StandardAuthorizer implements ClusterMetadataAuthorizer { */ private final CompletableFuture initialLoadFuture = new CompletableFuture<>(); + private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); + /** - * The current data. Can be read without a lock. Must be written while holding the object lock. + * The current data. We use a read-write lock to synchronize reads and writes to the data. */ private volatile StandardAuthorizerData data = StandardAuthorizerData.createEmpty(); + private void inWriteLock(Runnable function) { + lock.writeLock().lock(); + try { + function.run(); + } finally { + lock.writeLock().unlock(); + } + } + + private T inReadLock(Supplier function) { + lock.readLock().lock(); + try { + return function.get(); + } finally { + lock.readLock().unlock(); + } + } + @Override - public synchronized void setAclMutator(AclMutator aclMutator) { - this.data = data.copyWithNewAclMutator(aclMutator); + public void setAclMutator(AclMutator aclMutator) { + inWriteLock(() -> this.data = data.copyWithNewAclMutator(aclMutator)); } @Override public AclMutator aclMutatorOrException() { - AclMutator aclMutator = data.aclMutator; + AclMutator aclMutator = inReadLock(() -> data.aclMutator); if (aclMutator == null) { throw new NotControllerException("The current node is not the active controller."); } @@ -78,8 +100,8 @@ public AclMutator aclMutatorOrException() { } @Override - public synchronized void completeInitialLoad() { - data = data.copyWithNewLoadingComplete(true); + public void completeInitialLoad() { + inWriteLock(() -> data = data.copyWithNewLoadingComplete(true)); data.log.info("Completed initial ACL load process."); initialLoadFuture.complete(null); } @@ -97,17 +119,17 @@ public void completeInitialLoad(Exception e) { @Override public void addAcl(Uuid id, StandardAcl acl) { - data.addAcl(id, acl); + inWriteLock(() -> data.addAcl(id, acl)); } @Override public void removeAcl(Uuid id) { - data.removeAcl(id); + inWriteLock(() -> data.removeAcl(id)); } @Override - public synchronized void loadSnapshot(Map acls) { - data = data.copyWithNewAcls(acls.entrySet()); + public void loadSnapshot(Map acls) { + inWriteLock(() -> data = data.copyWithNewAcls(acls.entrySet())); } @Override @@ -129,23 +151,24 @@ public synchronized void loadSnapshot(Map acls) { public List authorize( AuthorizableRequestContext requestContext, List actions) { - StandardAuthorizerData curData = data; - List results = new ArrayList<>(actions.size()); - for (Action action: actions) { - AuthorizationResult result = curData.authorize(requestContext, action); - results.add(result); - } - return results; + return inReadLock(() -> { + List results = new ArrayList<>(actions.size()); + for (Action action : actions) { + AuthorizationResult result = data.authorize(requestContext, action); + results.add(result); + } + return results; + }); } @Override public Iterable acls(AclBindingFilter filter) { - return data.acls(filter); + return inReadLock(() -> data.acls(filter)); } @Override public int aclCount() { - return data.aclCount(); + return inReadLock(() -> data.aclCount()); } @Override @@ -156,7 +179,7 @@ public void close() throws IOException { } @Override - public synchronized void configure(Map configs) { + public void configure(Map configs) { Set superUsers = getConfiguredSuperUsers(configs); AuthorizationResult defaultResult = getDefaultResult(configs); int nodeId; @@ -165,17 +188,20 @@ public synchronized void configure(Map configs) { } catch (Exception e) { nodeId = -1; } - this.data = data.copyWithNewConfig(nodeId, superUsers, defaultResult); - this.data.log.info("set super.users={}, default result={}", String.join(",", superUsers), defaultResult); + final int finalNodeId = nodeId; + inWriteLock(() -> { + this.data = data.copyWithNewConfig(finalNodeId, superUsers, defaultResult); + this.data.log.info("set super.users={}, default result={}", String.join(",", superUsers), defaultResult); + }); } // VisibleForTesting Set superUsers() { - return new HashSet<>(data.superUsers()); + return inReadLock(() -> new HashSet<>(data.superUsers())); } AuthorizationResult defaultResult() { - return data.defaultResult(); + return inReadLock(() -> data.defaultResult()); } static Set getConfiguredSuperUsers(Map configs) { diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index d9ffd17562f53..a7aacea0a480e 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -39,13 +39,15 @@ import java.util.Collection; import java.util.Collections; import java.util.EnumSet; +import java.util.HashMap; import java.util.Iterator; +import java.util.List; import java.util.Map.Entry; import java.util.NavigableSet; import java.util.NoSuchElementException; import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentSkipListSet; +import java.util.TreeSet; +import java.util.stream.Collectors; import static org.apache.kafka.common.acl.AclOperation.ALL; import static org.apache.kafka.common.acl.AclOperation.ALTER; @@ -64,7 +66,7 @@ /** * A class which encapsulates the configuration and the ACL data owned by StandardAuthorizer. * - * The methods in this class support lockless concurrent access. + * The class is not thread-safe. */ public class StandardAuthorizerData { /** @@ -111,12 +113,12 @@ public class StandardAuthorizerData { /** * Contains all of the current ACLs sorted by (resource type, resource name). */ - private final ConcurrentSkipListSet aclsByResource; + private final TreeSet aclsByResource; /** * Contains all of the current ACLs indexed by UUID. */ - private final ConcurrentHashMap aclsById; + private final HashMap aclsById; private static Logger createLogger(int nodeId) { return new LogContext("[StandardAuthorizer " + nodeId + "] ").logger(StandardAuthorizerData.class); @@ -132,7 +134,7 @@ static StandardAuthorizerData createEmpty() { false, Collections.emptySet(), DENIED, - new ConcurrentSkipListSet<>(), new ConcurrentHashMap<>()); + new TreeSet<>(), new HashMap<>()); } private StandardAuthorizerData(Logger log, @@ -140,8 +142,8 @@ private StandardAuthorizerData(Logger log, boolean loadingComplete, Set superUsers, AuthorizationResult defaultResult, - ConcurrentSkipListSet aclsByResource, - ConcurrentHashMap aclsById) { + TreeSet aclsByResource, + HashMap aclsById) { this.log = log; this.auditLog = auditLogger(); this.aclMutator = aclMutator; @@ -193,8 +195,8 @@ StandardAuthorizerData copyWithNewAcls(Collection> aclE loadingComplete, superUsers, defaultRule.result, - new ConcurrentSkipListSet<>(), - new ConcurrentHashMap<>()); + new TreeSet<>(), + new HashMap<>()); for (Entry entry : aclEntries) { newData.addAcl(entry.getKey(), entry.getValue()); } @@ -534,49 +536,19 @@ Iterable acls(AclBindingFilter filter) { } class AclIterable implements Iterable { - private final AclBindingFilter filter; + private final List aclBindingList; AclIterable(AclBindingFilter filter) { - this.filter = filter; + this.aclBindingList = aclsByResource + .stream() + .map(StandardAcl::toBinding) + .filter(filter::matches) + .collect(Collectors.toList()); } @Override public Iterator iterator() { - return new AclIterator(filter); - } - } - - class AclIterator implements Iterator { - private final AclBindingFilter filter; - private final Iterator iterator; - private AclBinding next; - - AclIterator(AclBindingFilter filter) { - this.filter = filter; - this.iterator = aclsByResource.iterator(); - this.next = null; - } - - @Override - public boolean hasNext() { - while (next == null) { - if (!iterator.hasNext()) return false; - AclBinding binding = iterator.next().toBinding(); - if (filter.matches(binding)) { - next = binding; - } - } - return true; - } - - @Override - public AclBinding next() { - if (!hasNext()) { - throw new NoSuchElementException(); - } - AclBinding result = next; - next = null; - return result; + return aclBindingList.iterator(); } } From cc1bf8b39885c1853a6b13ad5ce428cbcdd87b11 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 12 Sep 2022 15:57:59 -0700 Subject: [PATCH 2/9] Benchmark changes --- .../kafka/jmh/acl/AclAuthorizerBenchmark.java | 64 +++++++++++++++---- .../authorizer/StandardAuthorizerData.java | 1 - 2 files changed, 53 insertions(+), 12 deletions(-) diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java index 65aa2a1f8d69b..9174c79157814 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java @@ -20,7 +20,9 @@ import kafka.security.authorizer.AclAuthorizer; import kafka.security.authorizer.AclAuthorizer.VersionedAcls; import kafka.security.authorizer.AclEntry; +import org.apache.kafka.common.Uuid; import org.apache.kafka.common.acl.AccessControlEntry; +import org.apache.kafka.common.acl.AclBinding; import org.apache.kafka.common.acl.AclBindingFilter; import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.acl.AclPermissionType; @@ -34,7 +36,10 @@ import org.apache.kafka.common.resource.ResourceType; import org.apache.kafka.common.security.auth.KafkaPrincipal; import org.apache.kafka.common.security.auth.SecurityProtocol; +import org.apache.kafka.metadata.authorizer.StandardAcl; +import org.apache.kafka.metadata.authorizer.StandardAuthorizer; import org.apache.kafka.server.authorizer.Action; +import org.apache.kafka.server.authorizer.Authorizer; import org.openjdk.jmh.annotations.Benchmark; import org.openjdk.jmh.annotations.BenchmarkMode; import org.openjdk.jmh.annotations.Fork; @@ -50,6 +55,7 @@ import org.openjdk.jmh.annotations.Warmup; import scala.collection.JavaConverters; +import java.io.IOException; import java.net.InetAddress; import java.util.ArrayList; import java.util.Collections; @@ -61,6 +67,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; @State(Scope.Benchmark) @Fork(value = 1) @@ -69,6 +76,22 @@ @BenchmarkMode(Mode.AverageTime) @OutputTimeUnit(TimeUnit.MILLISECONDS) public class AclAuthorizerBenchmark { + + public enum AuthorizerType { + ACL(AclAuthorizer::new), + KRAFT(StandardAuthorizer::new); + + private Supplier supplier; + + AuthorizerType(Supplier supplier) { + this.supplier = supplier; + } + + Authorizer newAuthorizer() { + return supplier.get(); + } + } + @Param({"10000", "50000", "200000"}) private int resourceCount; //no. of. rules per resource @@ -78,10 +101,13 @@ public class AclAuthorizerBenchmark { @Param({"0", "20", "50", "90", "99", "99.9", "99.99", "100"}) private double denyPercentage; + @Param({"ACL", "KRAFT"}) + private AuthorizerType authorizerType; + private final int hostPreCount = 1000; private final String resourceNamePrefix = "foo-bar35_resource-"; - private final AclAuthorizer aclAuthorizer = new AclAuthorizer(); private final KafkaPrincipal principal = new KafkaPrincipal(KafkaPrincipal.USER_TYPE, "test-user"); + private Authorizer authorizer; private List actions = new ArrayList<>(); private RequestContext authorizeContext; private RequestContext authorizeByResourceTypeContext; @@ -94,6 +120,7 @@ public class AclAuthorizerBenchmark { @Setup(Level.Trial) public void setup() throws Exception { + authorizer = authorizerType.newAuthorizer(); prepareAclCache(); prepareAclToUpdate(); // By adding `-95` to the resource name prefix, we cause the `TreeMap.from/to` call to return @@ -177,9 +204,23 @@ private void prepareAclCache() { } } + setupAcls(aclEntries); + } + + private void setupAcls(Map> aclEntries) { for (Map.Entry> entryMap : aclEntries.entrySet()) { - aclAuthorizer.updateCache(entryMap.getKey(), - new VersionedAcls(JavaConverters.asScalaSetConverter(entryMap.getValue()).asScala().toSet(), 1)); + switch (authorizerType) { + case ACL: + ((AclAuthorizer) authorizer).updateCache(entryMap.getKey(), + new VersionedAcls(JavaConverters.asScalaSetConverter(entryMap.getValue()).asScala().toSet(), 1)); + break; + case KRAFT: + for (AclEntry aclEntry : entryMap.getValue()) { + StandardAcl acl = StandardAcl.fromAclBinding(new AclBinding(entryMap.getKey(), aclEntry.ace())); + ((StandardAuthorizer) authorizer).addAcl(Uuid.randomUuid(), acl); + } + break; + } } } @@ -207,30 +248,31 @@ private Boolean shouldDeny() { } @TearDown(Level.Trial) - public void tearDown() { - aclAuthorizer.close(); + public void tearDown() throws IOException { + authorizer.close(); } @Benchmark public void testAclsIterator() { - aclAuthorizer.acls(AclBindingFilter.ANY); + authorizer.acls(AclBindingFilter.ANY); } @Benchmark public void testAuthorizer() { - aclAuthorizer.authorize(authorizeContext, actions); + authorizer.authorize(authorizeContext, actions); } @Benchmark public void testAuthorizeByResourceType() { - aclAuthorizer.authorizeByResourceType(authorizeByResourceTypeContext, AclOperation.READ, ResourceType.TOPIC); + authorizer.authorizeByResourceType(authorizeByResourceTypeContext, AclOperation.READ, ResourceType.TOPIC); } @Benchmark public void testUpdateCache() { - AclAuthorizer aclAuthorizer = new AclAuthorizer(); - for (Map.Entry e : aclToUpdate.entrySet()) { - aclAuthorizer.updateCache(e.getKey(), e.getValue()); + if (authorizerType == AuthorizerType.ACL) { + for (Map.Entry e : aclToUpdate.entrySet()) { + ((AclAuthorizer) authorizer).updateCache(e.getKey(), e.getValue()); + } } } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index a7aacea0a480e..f0d540e254298 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -44,7 +44,6 @@ import java.util.List; import java.util.Map.Entry; import java.util.NavigableSet; -import java.util.NoSuchElementException; import java.util.Set; import java.util.TreeSet; import java.util.stream.Collectors; From 3193034d4a2d032835c46b8059cae3cc3cf4d1f5 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 12 Sep 2022 16:55:50 -0700 Subject: [PATCH 3/9] Addressed the performance issues with StandardAuthorizer#acls --- .../kafka/jmh/acl/AclAuthorizerBenchmark.java | 1 + .../authorizer/StandardAuthorizerData.java | 14 ++++++++------ 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java index 9174c79157814..3a940cd1e6630 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java @@ -219,6 +219,7 @@ private void setupAcls(Map> aclEntries) { StandardAcl acl = StandardAcl.fromAclBinding(new AclBinding(entryMap.getKey(), aclEntry.ace())); ((StandardAuthorizer) authorizer).addAcl(Uuid.randomUuid(), acl); } + ((StandardAuthorizer) authorizer).completeInitialLoad(); break; } } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index f0d540e254298..2da55142cc9b7 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -36,6 +36,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.EnumSet; @@ -535,14 +536,15 @@ Iterable acls(AclBindingFilter filter) { } class AclIterable implements Iterable { - private final List aclBindingList; + private final List aclBindingList = new ArrayList<>(); AclIterable(AclBindingFilter filter) { - this.aclBindingList = aclsByResource - .stream() - .map(StandardAcl::toBinding) - .filter(filter::matches) - .collect(Collectors.toList()); + aclsByResource.forEach(acl -> { + AclBinding aclBinding = acl.toBinding(); + if (filter.matches(aclBinding)) { + aclBindingList.add(aclBinding); + } + }); } @Override From 70061c128a9c6c9652e88fbe8301f384f72d72de Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 12 Sep 2022 17:11:54 -0700 Subject: [PATCH 4/9] remove unused imports --- .../apache/kafka/metadata/authorizer/StandardAuthorizerData.java | 1 - 1 file changed, 1 deletion(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index 2da55142cc9b7..7842d27530361 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -47,7 +47,6 @@ import java.util.NavigableSet; import java.util.Set; import java.util.TreeSet; -import java.util.stream.Collectors; import static org.apache.kafka.common.acl.AclOperation.ALL; import static org.apache.kafka.common.acl.AclOperation.ALTER; From 6e0bdd0fc9e9938d728f340db60e38b2389f1fe7 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 12 Sep 2022 18:49:05 -0700 Subject: [PATCH 5/9] Addressed Jason's comments --- .../authorizer/StandardAuthorizer.java | 100 ++++++++++++------ .../authorizer/StandardAuthorizerData.java | 27 ++--- 2 files changed, 77 insertions(+), 50 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java index d8492590d5e58..ecf0119bee035 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java @@ -67,32 +67,25 @@ public class StandardAuthorizer implements ClusterMetadataAuthorizer { */ private volatile StandardAuthorizerData data = StandardAuthorizerData.createEmpty(); - private void inWriteLock(Runnable function) { + @Override + public void setAclMutator(AclMutator aclMutator) { lock.writeLock().lock(); try { - function.run(); + this.data = data.copyWithNewAclMutator(aclMutator); } finally { lock.writeLock().unlock(); } } - private T inReadLock(Supplier function) { + @Override + public AclMutator aclMutatorOrException() { + AclMutator aclMutator; lock.readLock().lock(); try { - return function.get(); + aclMutator = data.aclMutator; } finally { lock.readLock().unlock(); } - } - - @Override - public void setAclMutator(AclMutator aclMutator) { - inWriteLock(() -> this.data = data.copyWithNewAclMutator(aclMutator)); - } - - @Override - public AclMutator aclMutatorOrException() { - AclMutator aclMutator = inReadLock(() -> data.aclMutator); if (aclMutator == null) { throw new NotControllerException("The current node is not the active controller."); } @@ -101,7 +94,12 @@ public AclMutator aclMutatorOrException() { @Override public void completeInitialLoad() { - inWriteLock(() -> data = data.copyWithNewLoadingComplete(true)); + lock.writeLock().lock(); + try { + data = data.copyWithNewLoadingComplete(true); + } finally { + lock.writeLock().unlock(); + } data.log.info("Completed initial ACL load process."); initialLoadFuture.complete(null); } @@ -119,17 +117,32 @@ public void completeInitialLoad(Exception e) { @Override public void addAcl(Uuid id, StandardAcl acl) { - inWriteLock(() -> data.addAcl(id, acl)); + lock.writeLock().lock(); + try { + data.addAcl(id, acl); + } finally { + lock.writeLock().unlock(); + } } @Override public void removeAcl(Uuid id) { - inWriteLock(() -> data.removeAcl(id)); + lock.writeLock().lock(); + try { + data.removeAcl(id); + } finally { + lock.writeLock().unlock(); + } } @Override public void loadSnapshot(Map acls) { - inWriteLock(() -> data = data.copyWithNewAcls(acls.entrySet())); + lock.writeLock().lock(); + try { + data = data.copyWithNewAcls(acls.entrySet()); + } finally { + lock.writeLock().unlock(); + } } @Override @@ -151,24 +164,37 @@ public void loadSnapshot(Map acls) { public List authorize( AuthorizableRequestContext requestContext, List actions) { - return inReadLock(() -> { - List results = new ArrayList<>(actions.size()); + List results = new ArrayList<>(actions.size()); + lock.readLock().lock(); + try { for (Action action : actions) { AuthorizationResult result = data.authorize(requestContext, action); results.add(result); } - return results; - }); + } finally { + lock.readLock().unlock(); + } + return results; } @Override public Iterable acls(AclBindingFilter filter) { - return inReadLock(() -> data.acls(filter)); + lock.readLock().lock(); + try { + return data.acls(filter); + } finally { + lock.readLock().unlock(); + } } @Override public int aclCount() { - return inReadLock(() -> data.aclCount()); + lock.readLock().lock(); + try { + return data.aclCount(); + } finally { + lock.readLock().unlock(); + } } @Override @@ -188,20 +214,32 @@ public void configure(Map configs) { } catch (Exception e) { nodeId = -1; } - final int finalNodeId = nodeId; - inWriteLock(() -> { - this.data = data.copyWithNewConfig(finalNodeId, superUsers, defaultResult); - this.data.log.info("set super.users={}, default result={}", String.join(",", superUsers), defaultResult); - }); + lock.writeLock().lock(); + try { + data = data.copyWithNewConfig(nodeId, superUsers, defaultResult); + } finally { + lock.writeLock().unlock(); + } + this.data.log.info("set super.users={}, default result={}", String.join(",", superUsers), defaultResult); } // VisibleForTesting Set superUsers() { - return inReadLock(() -> new HashSet<>(data.superUsers())); + lock.readLock().lock(); + try { + return new HashSet<>(data.superUsers()); + } finally { + lock.readLock().unlock(); + } } AuthorizationResult defaultResult() { - return inReadLock(() -> data.defaultResult()); + lock.readLock().lock(); + try { + return data.defaultResult(); + } finally { + lock.readLock().unlock(); + } } static Set getConfiguredSuperUsers(Map configs) { diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index 7842d27530361..467911ae2f634 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -531,25 +531,14 @@ static AuthorizationResult findResult(Action action, } Iterable acls(AclBindingFilter filter) { - return new AclIterable(filter); - } - - class AclIterable implements Iterable { - private final List aclBindingList = new ArrayList<>(); - - AclIterable(AclBindingFilter filter) { - aclsByResource.forEach(acl -> { - AclBinding aclBinding = acl.toBinding(); - if (filter.matches(aclBinding)) { - aclBindingList.add(aclBinding); - } - }); - } - - @Override - public Iterator iterator() { - return aclBindingList.iterator(); - } + List aclBindingList = new ArrayList<>(); + aclsByResource.forEach(acl -> { + AclBinding aclBinding = acl.toBinding(); + if (filter.matches(aclBinding)) { + aclBindingList.add(aclBinding); + } + }); + return aclBindingList; } private interface MatchingRule { From b65180ddfe1cc06b35be2489282c7634b9fcb095 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Tue, 13 Sep 2022 09:10:46 -0700 Subject: [PATCH 6/9] clean imports --- .../org/apache/kafka/metadata/authorizer/StandardAuthorizer.java | 1 - 1 file changed, 1 deletion(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java index ecf0119bee035..b697068997dba 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java @@ -40,7 +40,6 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.locks.ReentrantReadWriteLock; -import java.util.function.Supplier; import static org.apache.kafka.server.authorizer.AuthorizationResult.ALLOWED; import static org.apache.kafka.server.authorizer.AuthorizationResult.DENIED; From 426e1a5e9814f4c4b289d8a6cc5d71f9481d00de Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 19 Sep 2022 19:17:16 -0700 Subject: [PATCH 7/9] revert benchmark changes --- .../kafka/jmh/acl/AclAuthorizerBenchmark.java | 65 ++++--------------- 1 file changed, 11 insertions(+), 54 deletions(-) diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java index 3a940cd1e6630..65aa2a1f8d69b 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/acl/AclAuthorizerBenchmark.java @@ -20,9 +20,7 @@ import kafka.security.authorizer.AclAuthorizer; import kafka.security.authorizer.AclAuthorizer.VersionedAcls; import kafka.security.authorizer.AclEntry; -import org.apache.kafka.common.Uuid; import org.apache.kafka.common.acl.AccessControlEntry; -import org.apache.kafka.common.acl.AclBinding; import org.apache.kafka.common.acl.AclBindingFilter; import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.acl.AclPermissionType; @@ -36,10 +34,7 @@ import org.apache.kafka.common.resource.ResourceType; import org.apache.kafka.common.security.auth.KafkaPrincipal; import org.apache.kafka.common.security.auth.SecurityProtocol; -import org.apache.kafka.metadata.authorizer.StandardAcl; -import org.apache.kafka.metadata.authorizer.StandardAuthorizer; import org.apache.kafka.server.authorizer.Action; -import org.apache.kafka.server.authorizer.Authorizer; import org.openjdk.jmh.annotations.Benchmark; import org.openjdk.jmh.annotations.BenchmarkMode; import org.openjdk.jmh.annotations.Fork; @@ -55,7 +50,6 @@ import org.openjdk.jmh.annotations.Warmup; import scala.collection.JavaConverters; -import java.io.IOException; import java.net.InetAddress; import java.util.ArrayList; import java.util.Collections; @@ -67,7 +61,6 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.TimeUnit; -import java.util.function.Supplier; @State(Scope.Benchmark) @Fork(value = 1) @@ -76,22 +69,6 @@ @BenchmarkMode(Mode.AverageTime) @OutputTimeUnit(TimeUnit.MILLISECONDS) public class AclAuthorizerBenchmark { - - public enum AuthorizerType { - ACL(AclAuthorizer::new), - KRAFT(StandardAuthorizer::new); - - private Supplier supplier; - - AuthorizerType(Supplier supplier) { - this.supplier = supplier; - } - - Authorizer newAuthorizer() { - return supplier.get(); - } - } - @Param({"10000", "50000", "200000"}) private int resourceCount; //no. of. rules per resource @@ -101,13 +78,10 @@ Authorizer newAuthorizer() { @Param({"0", "20", "50", "90", "99", "99.9", "99.99", "100"}) private double denyPercentage; - @Param({"ACL", "KRAFT"}) - private AuthorizerType authorizerType; - private final int hostPreCount = 1000; private final String resourceNamePrefix = "foo-bar35_resource-"; + private final AclAuthorizer aclAuthorizer = new AclAuthorizer(); private final KafkaPrincipal principal = new KafkaPrincipal(KafkaPrincipal.USER_TYPE, "test-user"); - private Authorizer authorizer; private List actions = new ArrayList<>(); private RequestContext authorizeContext; private RequestContext authorizeByResourceTypeContext; @@ -120,7 +94,6 @@ Authorizer newAuthorizer() { @Setup(Level.Trial) public void setup() throws Exception { - authorizer = authorizerType.newAuthorizer(); prepareAclCache(); prepareAclToUpdate(); // By adding `-95` to the resource name prefix, we cause the `TreeMap.from/to` call to return @@ -204,24 +177,9 @@ private void prepareAclCache() { } } - setupAcls(aclEntries); - } - - private void setupAcls(Map> aclEntries) { for (Map.Entry> entryMap : aclEntries.entrySet()) { - switch (authorizerType) { - case ACL: - ((AclAuthorizer) authorizer).updateCache(entryMap.getKey(), - new VersionedAcls(JavaConverters.asScalaSetConverter(entryMap.getValue()).asScala().toSet(), 1)); - break; - case KRAFT: - for (AclEntry aclEntry : entryMap.getValue()) { - StandardAcl acl = StandardAcl.fromAclBinding(new AclBinding(entryMap.getKey(), aclEntry.ace())); - ((StandardAuthorizer) authorizer).addAcl(Uuid.randomUuid(), acl); - } - ((StandardAuthorizer) authorizer).completeInitialLoad(); - break; - } + aclAuthorizer.updateCache(entryMap.getKey(), + new VersionedAcls(JavaConverters.asScalaSetConverter(entryMap.getValue()).asScala().toSet(), 1)); } } @@ -249,31 +207,30 @@ private Boolean shouldDeny() { } @TearDown(Level.Trial) - public void tearDown() throws IOException { - authorizer.close(); + public void tearDown() { + aclAuthorizer.close(); } @Benchmark public void testAclsIterator() { - authorizer.acls(AclBindingFilter.ANY); + aclAuthorizer.acls(AclBindingFilter.ANY); } @Benchmark public void testAuthorizer() { - authorizer.authorize(authorizeContext, actions); + aclAuthorizer.authorize(authorizeContext, actions); } @Benchmark public void testAuthorizeByResourceType() { - authorizer.authorizeByResourceType(authorizeByResourceTypeContext, AclOperation.READ, ResourceType.TOPIC); + aclAuthorizer.authorizeByResourceType(authorizeByResourceTypeContext, AclOperation.READ, ResourceType.TOPIC); } @Benchmark public void testUpdateCache() { - if (authorizerType == AuthorizerType.ACL) { - for (Map.Entry e : aclToUpdate.entrySet()) { - ((AclAuthorizer) authorizer).updateCache(e.getKey(), e.getValue()); - } + AclAuthorizer aclAuthorizer = new AclAuthorizer(); + for (Map.Entry e : aclToUpdate.entrySet()) { + aclAuthorizer.updateCache(e.getKey(), e.getValue()); } } } From 3171facfbfc4085efaded9cb4e854eaa5c26b6f8 Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Mon, 19 Sep 2022 21:52:32 -0700 Subject: [PATCH 8/9] reduce the crictical section for write lock on loading snapshot --- .../authorizer/StandardAuthorizer.java | 6 ++++- .../authorizer/StandardAuthorizerData.java | 22 +++++++++++++++++++ 2 files changed, 27 insertions(+), 1 deletion(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java index b697068997dba..799c80a9a0516 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java @@ -136,9 +136,13 @@ public void removeAcl(Uuid id) { @Override public void loadSnapshot(Map acls) { + StandardAuthorizerData newData = StandardAuthorizerData.createEmpty(); + for (Map.Entry entry : acls.entrySet()) { + newData.addAcl(entry.getKey(), entry.getValue()); + } lock.writeLock().lock(); try { - data = data.copyWithNewAcls(acls.entrySet()); + data = data.copyWithNewAcls(newData.getAclsByResource(), newData.getAclsById()); } finally { lock.writeLock().unlock(); } diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index 467911ae2f634..d90de070e220d 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -203,6 +203,20 @@ StandardAuthorizerData copyWithNewAcls(Collection> aclE return newData; } + StandardAuthorizerData copyWithNewAcls(TreeSet aclsByResource, HashMap aclsById) { + StandardAuthorizerData newData = new StandardAuthorizerData( + log, + aclMutator, + loadingComplete, + superUsers, + defaultRule.result, + aclsByResource, + aclsById); + log.info("Initialized with {} acl(s).", aclsById.size()); + return newData; + } + void addAcl(Uuid id, StandardAcl acl) { try { StandardAcl prevAcl = aclsById.putIfAbsent(id, acl); @@ -615,4 +629,12 @@ MatchingAclRule build() { } } } + + TreeSet getAclsByResource() { + return aclsByResource; + } + + HashMap getAclsById() { + return aclsById; + } } From 442df1285c4f0739d4adaceac99244bde400728b Mon Sep 17 00:00:00 2001 From: Akhilesh Chaganti Date: Tue, 20 Sep 2022 10:19:13 -0700 Subject: [PATCH 9/9] Address David's comments --- .../authorizer/StandardAuthorizer.java | 9 +++++-- .../authorizer/StandardAuthorizerData.java | 24 +++++-------------- 2 files changed, 13 insertions(+), 20 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java index 799c80a9a0516..197272a3e66f5 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizer.java @@ -62,7 +62,9 @@ public class StandardAuthorizer implements ClusterMetadataAuthorizer { private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock(); /** - * The current data. We use a read-write lock to synchronize reads and writes to the data. + * The current data. We use a read-write lock to synchronize reads and writes to the data. We + * expect one writer and multiple readers accessing the ACL data, and we use the lock to make + * sure we have consistent reads when writer tries to change the data. */ private volatile StandardAuthorizerData data = StandardAuthorizerData.createEmpty(); @@ -170,8 +172,9 @@ public List authorize( List results = new ArrayList<>(actions.size()); lock.readLock().lock(); try { + StandardAuthorizerData curData = data; for (Action action : actions) { - AuthorizationResult result = data.authorize(requestContext, action); + AuthorizationResult result = curData.authorize(requestContext, action); results.add(result); } } finally { @@ -184,6 +187,8 @@ public List authorize( public Iterable acls(AclBindingFilter filter) { lock.readLock().lock(); try { + // The Iterable returned here is consistent because it is created over a read-only + // copy of ACLs data. return data.acls(filter); } finally { lock.readLock().unlock(); diff --git a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java index d90de070e220d..c6e3b74a2ab03 100644 --- a/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java +++ b/metadata/src/main/java/org/apache/kafka/metadata/authorizer/StandardAuthorizerData.java @@ -37,13 +37,11 @@ import org.slf4j.LoggerFactory; import java.util.ArrayList; -import java.util.Collection; import java.util.Collections; import java.util.EnumSet; import java.util.HashMap; import java.util.Iterator; import java.util.List; -import java.util.Map.Entry; import java.util.NavigableSet; import java.util.Set; import java.util.TreeSet; @@ -187,22 +185,6 @@ StandardAuthorizerData copyWithNewConfig(int nodeId, aclsById); } - StandardAuthorizerData copyWithNewAcls(Collection> aclEntries) { - StandardAuthorizerData newData = new StandardAuthorizerData( - log, - aclMutator, - loadingComplete, - superUsers, - defaultRule.result, - new TreeSet<>(), - new HashMap<>()); - for (Entry entry : aclEntries) { - newData.addAcl(entry.getKey(), entry.getValue()); - } - log.info("Applied {} acl(s) from image.", aclEntries.size()); - return newData; - } - StandardAuthorizerData copyWithNewAcls(TreeSet aclsByResource, HashMap aclsById) { StandardAuthorizerData newData = new StandardAuthorizerData( @@ -544,6 +526,12 @@ static AuthorizationResult findResult(Action action, return acl.permissionType().equals(ALLOW) ? ALLOWED : DENIED; } + /** + * Creates a consistent Iterable on read-only copy of AclBindings data for the given filter. + * + * @param filter The filter constraining the AclBindings to be present in the Iterable. + * @return Iterable over AclBindings matching the filter. + */ Iterable acls(AclBindingFilter filter) { List aclBindingList = new ArrayList<>(); aclsByResource.forEach(acl -> {