如何在Spring Boot项目中配置Kafka ACL及相关操作(用户创建、读写权限等)
Got it, let's walk through how to set up Kafka ACLs and manage users/permissions in your Spring Boot project. I’ve tackled this scenario multiple times, so here’s a practical, step-by-step guide that should cover everything you need:
Before configuring Spring Boot, your Kafka server needs to be set up to use ACLs. Update your server.properties file with these key settings:
# Enable the built-in ACL authorizer authorizer.class.name=kafka.security.authorizer.AclAuthorizer # Define a super user (this account will have full access to manage ACLs) super.users=User:admin # Set listener protocol (use SSL in production, PLAINTEXT for testing) listeners=PLAINTEXT://localhost:9092 security.inter.broker.protocol=PLAINTEXT
Restart your Kafka brokers after making these changes.
Next, configure your Spring Boot app to authenticate with Kafka and respect ACLs. Add these settings to your application.yml (or application.properties):
spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer consumer: key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer group-id: my-app-consumer-group properties: security.protocol: PLAINTEXT sasl.mechanism: PLAIN # JAAS config for authenticating as the super user (admin) to manage ACLs sasl.jaas.config: org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secure-password";
Pro tip: In production, never hardcode passwords—use environment variables or Spring Cloud Config to inject sensitive values.
Spring Boot can auto-configure a AdminClient bean that lets you interact with Kafka's admin APIs. Here's how to use it to create users (bind ACLs to them) and manage permissions:
Step 3.1: Configure the AdminClient Bean
Create a configuration class to expose the AdminClient with your security settings:
import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaAdminConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; @Value("${spring.kafka.properties.sasl.jaas.config}") private String jaasConfig; @Bean public AdminClient adminClient() { Map<String, Object> config = new HashMap<>(); config.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); config.put("security.protocol", "PLAINTEXT"); config.put("sasl.mechanism", "PLAIN"); config.put("sasl.jaas.config", jaasConfig); return AdminClient.create(config); } }
Step 3.2: Create a Service for ACL Operations
Build a service class to encapsulate user and ACL management logic:
import org.apache.kafka.clients.admin.*; import org.apache.kafka.common.acl.*; import org.apache.kafka.common.resource.ResourcePattern; import org.apache.kafka.common.resource.ResourcePatternType; import org.apache.kafka.common.resource.ResourceType; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.util.Collections; import java.util.concurrent.ExecutionException; @Service public class KafkaAclManager { private final AdminClient adminClient; @Autowired public KafkaAclManager(AdminClient adminClient) { this.adminClient = adminClient; } // Grant produce + consume permissions to a user for a specific topic public void grantTopicPermissions(String username, String topicName, String consumerGroup) throws ExecutionException, InterruptedException { // 1. Allow writing to the topic (producer access) AclBinding producerAcl = new AclBinding( new ResourcePattern(ResourceType.TOPIC, topicName, ResourcePatternType.LITERAL), new AccessControlEntry( "User:" + username, "*", // Allow from any host AclOperation.WRITE, AclPermissionType.ALLOW ) ); // 2. Allow reading from the topic (consumer access) AclBinding consumerTopicAcl = new AclBinding( new ResourcePattern(ResourceType.TOPIC, topicName, ResourcePatternType.LITERAL), new AccessControlEntry( "User:" + username, "*", AclOperation.READ, AclPermissionType.ALLOW ) ); // 3. Allow accessing the consumer group (critical for consumers) AclBinding consumerGroupAcl = new AclBinding( new ResourcePattern(ResourceType.GROUP, consumerGroup, ResourcePatternType.LITERAL), new AccessControlEntry( "User:" + username, "*", AclOperation.READ, AclPermissionType.ALLOW ) ); // Batch create all ACLs CreateAclsResult result = adminClient.createAcls( Collections.singletonList(producerAcl, consumerTopicAcl, consumerGroupAcl) ); result.all().get(); // Wait for operation to complete System.out.println("Permissions granted to user: " + username); } // List all ACLs for a specific user public void listUserAcls(String username) throws ExecutionException, InterruptedException { AclBindingFilter filter = new AclBindingFilter( null, new AccessControlEntryFilter("User:" + username, "*", null, null) ); ListAclsResult result = adminClient.listAcls(filter); result.values().get().forEach(acl -> System.out.println("User ACL: " + acl) ); } // Revoke all ACLs for a user public void revokeAllUserPermissions(String username) throws ExecutionException, InterruptedException { AclBindingFilter filter = new AclBindingFilter( null, new AccessControlEntryFilter("User:" + username, "*", null, null) ); DeleteAclsResult result = adminClient.deleteAcls(Collections.singletonList(filter)); result.all().get(); System.out.println("All permissions revoked for user: " + username); } }
Note: Kafka doesn't have a "create user" API directly—users are tied to your authentication system (like JAAS, LDAP, or OAuth). Here, we're binding ACLs to a username that should exist in your auth setup.
You can expose these operations via a REST controller or test them in a service:
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.*; import java.util.concurrent.ExecutionException; @RestController @RequestMapping("/kafka-acl") public class KafkaAclController { @Autowired private KafkaAclManager aclManager; @PostMapping("/grant-permissions") public String grantPermissions( @RequestParam String username, @RequestParam String topicName, @RequestParam String consumerGroup) { try { aclManager.grantTopicPermissions(username, topicName, consumerGroup); return "Permissions granted successfully for user: " + username; } catch (ExecutionException | InterruptedException e) { return "Error granting permissions: " + e.getMessage(); } } }
- SSL in Production: Replace
PLAINTEXTwithSSLand configure keystore/truststore paths for secure communication. - ACL Priority:
DENYACLs take precedence overALLOWACLs—be careful with overlapping rules. - Super User Lockdown: Restrict access to the super user account (admin) to only trusted services/people.
内容的提问来源于stack exchange,提问作者Nagihan Sevgi

