You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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:

1. First, Enable ACLs on Your Kafka Cluster

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.

2. Spring Boot Client Configuration

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.

3. Manage Users & ACLs via Spring Code

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.

4. Test the 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();
        }
    }
}
5. Key Things to Remember
  • SSL in Production: Replace PLAINTEXT with SSL and configure keystore/truststore paths for secure communication.
  • ACL Priority: DENY ACLs take precedence over ALLOW ACLs—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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 15:07:35