11.3 Authentication
KafkaPrincipal is established during authentication based on the protocol (e.g. User:Alice) and can be customized via principal.builder.class.
"Once authenticated, Alice's identity is associated with the connection throughout the LIFETIME of the connection. Kafka uses an instance of
KafkaPrincipalto represent client identity and uses this principal to grant access to resources and allocate quotas."
KafkaPrincipal is established during authentication based on the protocol (e.g. User:Alice) and can be customized via principal.builder.class.
ANONYMOUS CONNECTIONS
"The principal
User:ANONYMOUSis used for unauthenticated connections. This includes clients on PLAINTEXT listeners as well as unauthenticated clients on SSL listeners."
3.1 SSL authentication
"When a connection is established over TLS, the TLS handshake process performs authentication, negotiates cryptographic parameters, and generates shared keys for encryption. The server's digital certificate is verified by the client to establish the identity of the server. If client authentication using SSL is enabled, the server also verifies the client's digital certificate."
⚠️ SSL PERFORMANCE
"SSL channels are encrypted and hence introduce a noticeable overhead in terms of CPU usage. ZERO-COPY TRANSFER IS CURRENTLY NOT SUPPORTED FOR SSL. Depending on the traffic pattern, the overhead may be up to 20–30%."
(This is the same zero-copy loss from Ch. 6 §5.4 — and it's exactly why Ch. 10 §7.2 says to consider consuming locally and producing remotely when only the WAN hop needs encryption.)
What goes where
| Broker | Client |
|---|---|
Key store (required)
| Trust store (required)
|
| Trust store — only if TLS is used for inter-broker communication or client authentication is enabled | Key store — only if client authentication is required
|
"Broker certificates should contain the broker hostname as a Subject Alternative Name (SAN) extension or as the Common Name (CN) to enable clients to verify the server hostname. Wildcard certificates can be used to simplify administration by using the same key store for all brokers in a domain."
⚠️ SERVER HOSTNAME VERIFICATION
"By default, Kafka clients verify that the hostname of the server stored in the server certificate matches the host that the client is connecting to. The connection hostname may be a bootstrap server the client is configured with or an ADVERTISED LISTENER hostname returned by a broker in a metadata response. HOSTNAME VERIFICATION IS A CRITICAL PART OF SERVER AUTHENTICATION THAT PROTECTS AGAINST MAN-IN-THE-MIDDLE ATTACKS AND HENCE SHOULD NOT BE DISABLED IN PRODUCTION SYSTEMS."
(This is the flip side of Ch. 5 §3.1's DNS-alias problem: hostname verification is why a DNS alias breaks SASL, and disabling verification is the wrong fix.)
Client authentication modes
| Setting | Behaviour |
|---|---|
ssl.client.auth=required | Clients must present a valid certificate. |
ssl.client.auth=requested | “Clients that are not configured with key stores will complete the TLS handshake in this case, but will be assigned the principal User:ANONYMOUS.” |
SASL_SSL listeners | “Disable TLS client authentication and rely on SASL authentication and the KafkaPrincipal established by SASL.” |
"By default, the distinguished name (DN) of the client certificate is used as the
KafkaPrincipalfor authorization and quotas. The configuration optionssl.principal.mapping.rulescan be used to provide a list of rules to customize the principal."
⚠️ "If SSL is used for inter-broker communication, broker trust stores should include the CA of the BROKER certificates AS WELL AS the CA of the CLIENT certificates."
Certificate generation — the full sequence
Step 1 — self-signed CA for brokers:
# Create the CA key-pair; "We use this for signing certificates."
keytool -genkeypair -keyalg RSA -keysize 2048 -keystore server.ca.p12 \
-storetype PKCS12 -storepass server-ca-password -keypass server-ca-password \
-alias ca -dname "CN=BrokerCA" -ext bc=ca:true -validity 365
# Export the CA's public certificate; goes into trust stores + cert chains
keytool -export -file server.ca.crt -keystore server.ca.p12 \
-storetype PKCS12 -storepass server-ca-password -alias ca -rfcStep 2 — broker key store with a CA-signed certificate:
# ① generate the broker's private key
keytool -genkey -keyalg RSA -keysize 2048 -keystore server.ks.p12 \
-storepass server-ks-password -keypass server-ks-password -alias server \
-storetype PKCS12 -dname "CN=Kafka,O=Confluent,C=GB" -validity 365
# ② certificate signing request
keytool -certreq -file server.csr -keystore server.ks.p12 -storetype PKCS12 \
-storepass server-ks-password -keypass server-ks-password -alias server
# ③ sign it with the CA — NOTE THE SAN EXTENSION (hostname verification!)
keytool -gencert -infile server.csr -outfile server.crt \
-keystore server.ca.p12 -storetype PKCS12 -storepass server-ca-password \
-alias ca -ext SAN=DNS:broker1.example.com -validity 365
# ④ import the certificate CHAIN back into the broker key store
cat server.crt server.ca.crt > serverchain.crt
keytool -importcert -file serverchain.crt -keystore server.ks.p12 \
-storepass server-ks-password -keypass server-ks-password -alias server \
-storetype PKCS12 -noprompt"If using wildcard hostnames, the same key store can be used for all brokers. Otherwise, create a key store for EACH broker with its fully qualified domain name (FQDN)."
Step 3 — trust stores:
# broker trust store (for inter-broker TLS): contains the BROKER CA
keytool -import -file server.ca.crt -keystore server.ts.p12 \
-storetype PKCS12 -storepass server-ts-password -alias server -noprompt
# client trust store: contains the BROKER CA
keytool -import -file server.ca.crt -keystore client.ts.p12 \
-storetype PKCS12 -storepass client-ts-password -alias ca -nopromptStep 4 — client CA + client key store (only if ssl.client.auth is on):
# a SEPARATE CA for clients
keytool -genkeypair ... -keystore client.ca.p12 -dname CN=ClientCA -ext bc=ca:true
keytool -export -file client.ca.crt -keystore client.ca.p12 ...
# client key store; the DN becomes the principal:
# User:CN=Metrics App,O=Confluent,C=GB
keytool -genkey ... -keystore client.ks.p12 \
-dname "CN=Metrics App,O=Confluent,C=GB" -validity 365
keytool -certreq ... ; keytool -gencert ... ;
cat client.crt client.ca.crt > clientchain.crt
keytool -importcert -file clientchain.crt -keystore client.ks.p12 ...
# ⚠ "The broker's trust store should contain the CAs of ALL clients."
keytool -import -file client.ca.crt -keystore server.ts.p12 -alias client ...Broker and client TLS configuration
# BROKER
ssl.keystore.location=/path/to/server.ks.p12
ssl.keystore.password=server-ks-password
ssl.key.password=server-ks-password
ssl.keystore.type=PKCS12
ssl.truststore.location=/path/to/server.ts.p12
ssl.truststore.password=server-ts-password
ssl.truststore.type=PKCS12
ssl.client.auth=required# CLIENT
ssl.truststore.location=/path/to/client.ts.p12
ssl.truststore.password=client-ts-password
ssl.truststore.type=PKCS12
ssl.keystore.location=/path/to/client.ks.p12 # only if client auth
ssl.keystore.password=client-ks-password
ssl.key.password=client-ks-password
ssl.keystore.type=PKCS12TRUST STORES — you may not need them
"Trust store configuration can be omitted in brokers as well as clients when using certificates signed by WELL-KNOWN TRUSTED AUTHORITIES. The default trust stores in the Java installation will be sufficient to establish trust in this case."
💡 Rotating certificates without a restart
"Key stores and trust stores must be updated periodically BEFORE CERTIFICATES EXPIRE to avoid TLS handshake failures. Broker SSL stores can be DYNAMICALLY UPDATED by modifying the same file or setting the configuration option to a new versioned file. In both cases, the Admin API or the Kafka configs tool can be used to trigger the update."
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--command-config admin.props \
--entity-type brokers --entity-name 0 --alter --add-config \
'listener.name.external.ssl.keystore.location=/path/to/server.ks.p12'TLS security considerations
- ✓ “Kafka enables Only the newer protocols TLSv1.2 and TLSv1.3 by default, since older protocols like TLSv1.1 have Known vulnerabilities.”
- ✓ “Due to issues with Insecure renegotiation, Kafka Does not support renegotiation for TLS connections.”
- ✓ “Hostname verification is enabled by default to prevent MITM attacks.”
- ✓ Restrict cipher suites: “Strong ciphers with At least a 256-bit encryption key size protect against cryptographic attacks.” — some organizations must comply with FIPS 140-2.
- ⚠ “Since key stores containing Private keys are stored on the filesystem by default, It is vital to limit access to key store files using filesystem permissions.”
- ✓ “Standard Java TLS features can be used to enable Certificate revocation if a private key is compromised. Short-lived keys can be used to reduce exposure.”
- ⚠ DoS risk: “TLS handshakes are Expensive and utilize a significant amount of time On network threads in brokers. Listeners using TLS on insecure networks should be protected against DoS attacks using Connection quotas and limits.”
- ►
connection.failed.authentication.delay.ms— “delay failed response on authentication failures to Reduce the rate at which authentication failures are retried by clients.”
- ►
(The DoS point connects to Ch. 6 §5.1: TLS handshakes burn network thread time — the same finite pool that moves every request onto the request queue.)
3.2 SASL — the framework and four mechanisms
"SASL authentication is performed through a sequence of server challenges and client responses where the SASL mechanism defines the sequence and wire format of challenges and responses."
| Mechanism | What it is | Production readiness |
|---|---|---|
| GSSAPI | "Kerberos authentication... can be used to integrate with Kerberos servers like Active Directory or OpenLDAP" | Production-ready |
| PLAIN | "Username/password authentication that is typically used with a CUSTOM SERVER-SIDE CALLBACK to verify passwords from an external password store" | ⚠️ Built-in store is insecure |
| SCRAM-SHA-256 / SCRAM-SHA-512 | "Username/password authentication available out of the box WITHOUT the need for additional password stores" | ✅ if ZooKeeper is secure |
| OAUTHBEARER | "Authentication using OAuth bearer tokens, typically used with custom callbacks to acquire and validate tokens" | ⚠️ Built-in impl is not for production |
Enabling and selecting:
# broker, per listener
sasl.enabled.mechanisms=<one or more>
# client
sasl.mechanism=<one of the enabled>JAAS configuration:
"Kafka uses the Java Authentication and Authorization Service (JAAS) for configuring SASL. The configuration option
sasl.jaas.configcontains a single JAAS configuration entry that specifies a login module and its options. Brokers use the listener AND mechanism prefixes" — e.g.listener.name.external.gssapi.sasl.jaas.config.
JAAS CONFIGURATION FILE vs
sasl.jaas.config"JAAS configuration may also be specified in configuration files using the Java system property
java.security.auth.login.config. However, the Kafka optionsasl.jaas.configis RECOMMENDED since it supports PASSWORD PROTECTION and SEPARATE CONFIGURATION FOR EACH SASL MECHANISM when multiple mechanisms are enabled on a listener."
The three customization hooks — memorize what each is for:
- Login callback handler (brokers or clients) — “customize the login process, for example, to acquire credentials to be used for authentication”.
- Server callback handler (brokers) — “perform authentication of client credentials, for example, to verify passwords using an external password server”.
- Client callback handler (clients) — “inject client credentials instead of including them in the JAAS configuration”.
- SASL authentication itself is “a sequence of server challenges and client responses”; the mechanism defines the sequence and the wire format.
3.3 SASL/GSSAPI (Kerberos)
Broker configuration:
sasl.enabled.mechanisms=GSSAPI
listener.name.external.gssapi.sasl.jaas.config=\
com.sun.security.auth.module.Krb5LoginModule required \
useKeyTab=true storeKey=true \
keyTab="/path/to/broker1.keytab" \
principal="kafka/broker1.example.com@EXAMPLE.COM";- "Keytab files must be readable by the broker process."
- ⚠️ "Service principal for brokers SHOULD INCLUDE THE BROKER HOSTNAME" — "Broker hostnames are verified by clients to ensure server authenticity and prevent man-in-the-middle attacks."
Inter-broker:
sasl.mechanism.inter.broker.protocol=GSSAPI
sasl.kerberos.service.name=kafkaClient:
sasl.mechanism=GSSAPI
sasl.kerberos.service.name=kafka # the name of the service you connect TO
sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required \
useKeyTab=true storeKey=true \
keyTab="/path/to/alice.keytab" \
principal="Alice@EXAMPLE.COM"; # clients MAY omit the hostnamePrincipal derivation: "The SHORT NAME of the principal is used as the client identity by default. For example, User:Alice is the client principal and User:kafka is the broker principal." Customize with sasl.kerberos.principal.to.local.rules.
⚠️ DNS requirement: "Kerberos requires a SECURE DNS SERVICE for hostname lookup during authentication. In deployments where forward and reverse lookup DO NOT MATCH, the Kerberos configuration file
krb5.confon clients can be configured to setrdns=falseto disable reverse lookup."
Kerberos security considerations
- ✓ “Use of SASL_SSL is Recommended in production deployments using Kerberos”
- Why: “If TLS is not used, Eavesdroppers on the network may gain enough information to mount a dictionary attack or brute-force attack to steal client credentials.”
- ✓ “It is Safer to use Randomly generated keys for brokers instead of keys generated from Passwords that are easier to crack.”
- ✓ “Weak encryption algorithms like des-md5 should be avoided.”
- ⚠ “Access to Keytab files must be restricted using filesystem permissions Since any user in possession of the file may impersonate the user.”
- ⚠ Availability dependency: “Because DoS attacks against the KDC or DNS service Can result in authentication failures in clients, it is necessary to monitor the availability of these services.”
- ⚠ “Kerberos also relies on Loosely synchronized clocks with configurable variability To detect replay attacks. It is important to ensure that Clock synchronization is secure.”
That last one is subtle and worth flagging: NTP becomes part of your security perimeter. An attacker who can skew clocks can defeat replay detection.
3.4 SASL/PLAIN
The default (insecure) implementation uses the broker's JAAS config as the password store:
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
listener.name.external.plain.sasl.jaas.config=\
org.apache.kafka.common.security.plain.PlainLoginModule required \
username="kafka" password="kafka-password" \ # ① broker's own creds
user_kafka="kafka-password" \ # ② the password STORE
user_Alice="Alice-password";# CLIENT
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule \
required username="Alice" password="Alice-password";⚠️ "The built-in implementation that stores ALL passwords in EVERY broker's JAAS configuration is INSECURE AND NOT VERY FLEXIBLE since ALL BROKERS WILL NEED TO BE RESTARTED TO ADD OR REMOVE A USER."
The production pattern: a custom server callback handler
Two jobs: integrate with a secure third-party password server, and support password rotation.
💡 "On the server side, a server callback handler should support BOTH OLD AND NEW PASSWORDS for an OVERLAPPING PERIOD until all clients switch to the new password."
public class PasswordVerifier extends PlainServerCallbackHandler {
private final List<String> passwdFiles = new ArrayList<>(); // ① multiple
// files →
@Override // rotation
public void configure(Map<String, ?> configs, String mechanism,
List<AppConfigurationEntry> jaasEntries) {
Map<String,?> loginOptions = jaasEntries.get(0).getOptions();
String files = (String) loginOptions.get("password.files"); // ② JAAS option
Collections.addAll(passwdFiles, files.split(","));
}
@Override
protected boolean authenticate(String user, char[] password) {
return passwdFiles.stream() // ③ match ANY
.anyMatch(file -> authenticate(file, user, password)); // file
}
private boolean authenticate(String file, String user, char[] password) {
try {
String cmd = String.format("htpasswd -vb %s %s %s", // ④ htpasswd
file, user, new String(password)); // for
return Runtime.getRuntime().exec(cmd).waitFor() == 0; // simplicity
} catch (Exception e) {
return false;
}
}
}④ "We use htpasswd for simplicity. A secure database can be used for production deployments."
listener.name.external.plain.sasl.jaas.config=\
org.apache.kafka.common.security.plain.PlainLoginModule required \
password.files="/path/to/htpassword.props,/path/to/oldhtpassword.props";
listener.name.external.plain.sasl.server.callback.handler.class=\
com.example.PasswordVerifierClient-side callback: load passwords at connection time, not startup
"a client callback handler that implements
org.apache.kafka.common.security.auth.AuthenticateCallbackHandlercan be used to load passwords DYNAMICALLY AT RUNTIME when a connection is established instead of loading statically from the JAAS configuration during startup."
@Override
public void handle(Callback[] callbacks) throws IOException {
Properties props = Utils.loadProps(passwdFile); // ① reload EVERY time
PasswordConfig config = new PasswordConfig(props); // → supports rotation
String user = config.getString("username");
String password = config.getPassword("password").value();// ② returns the real
for (Callback callback: callbacks) { // value even if
if (callback instanceof NameCallback) // EXTERNALIZED
((NameCallback) callback).setName(user);
else if (callback instanceof PasswordCallback) {
((PasswordCallback) callback).setPassword(password.toCharArray());
}
}
}
private static class PasswordConfig extends AbstractConfig {
static ConfigDef CONFIG = new ConfigDef()
.define("username", STRING, HIGH, "User name")
.define("password", PASSWORD, HIGH, "User password"); // ③ PASSWORD type
PasswordConfig(Properties props) { super(CONFIG, props, false); }
}③ 💡 "We define password configs with the PASSWORD type to ensure that passwords are NOT INCLUDED IN LOG ENTRIES." — a cheap, high-value trick.
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule \
required file="/path/to/credentials.props";
sasl.client.callback.handler.class=com.example.PasswordProviderPLAIN security considerations
- ⚠ “Since SASL/PLAIN transmits Clear-text passwords over the wire, the PLAIN mechanism Should be enabled only with encryption using sasl_ssl.”
- ⚠ “Passwords stored in clear text in the JAAS configuration of brokers And clients are Not secure — consider Encrypting or externalizing.”
- ✓ “Use a Secure external password server that stores passwords securely and Enforces strong password policies.”
► Clear-text passwords: “Avoid clear-text passwords in configuration files Even if the files can be protected using filesystem permissions. Consider externalizing or encrypting passwords to ensure that passwords are Not inadvertently exposed.”
3.5 SASL/SCRAM
"RFC-5802 introduces a secure username/password authentication mechanism that addresses the security concerns with password authentication mechanisms like SASL/PLAIN, which send passwords over the wire. The Salted Challenge Response Authentication Mechanism (SCRAM) avoids transmitting clear-text passwords and stores passwords in a format that makes it IMPRACTICAL TO IMPERSONATE CLIENTS. Salting combines passwords with some random data before applying a one-way cryptographic hash function."
💡 “Kafka has a built-in SCRAM provider that can be used in deployments with secure ZooKeeper without the need for additional password servers.” SCRAM is the “batteries-included” secure option. That's its appeal.
Creating users — note this uses --zookeeper, before brokers start:
bin/kafka-configs.sh --zookeeper localhost:2181 --alter --add-config \
'SCRAM-SHA-512=[iterations=8192,password=Alice-password]' \
--entity-type users --entity-name Alice"An initial set of users can be created after starting ZooKeeper PRIOR TO STARTING BROKERS. Brokers load SCRAM user metadata into an IN-MEMORY CACHE during startup, ensuring that all users, including the broker user for inter-broker communication, can authenticate successfully. Users can be added or deleted AT ANY TIME. Brokers keep the cache up-to-date using notifications based on a ZOOKEEPER WATCHER."
# BROKER
sasl.enabled.mechanisms=SCRAM-SHA-512
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-512
listener.name.external.scram-sha-512.sasl.jaas.config=\
org.apache.kafka.common.security.scram.ScramLoginModule required \
username="kafka" password="kafka-password";# CLIENT
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule \
required username="Alice" password="Alice-password";Deleting a user:
bin/kafka-configs.sh --zookeeper localhost:2181 --alter --delete-config \
'SCRAM-SHA-512' --entity-type users --entity-name Alice⚠️ "When an existing user is deleted, new connections cannot be established for that user, BUT EXISTING CONNECTIONS OF THE USER WILL CONTINUE TO WORK. A REAUTHENTICATION INTERVAL can be configured for the broker to LIMIT THE AMOUNT OF TIME existing connections may continue to operate after a user is deleted."
SCRAM security considerations
- ✔ Kafka's safeguards: “supporting only the strong hashing algorithms SHA-256 and SHA-512 and avoiding weaker algorithms like SHA-1”, plus “a high default iteration count of 4,096 and unique random salts for every stored key to limit the impact if ZooKeeper security is compromised.”
- ⚠ “any password-based system is only as secure as the passwords. Strong password policies must be enforced.”
⚠ Required combination:
- “SCRAM must be used with
SASL_SSLto avoid eavesdroppers from gaining access to hashed keys during authentication.” - “ZooKeeper must also be SSL-enabled.”
- “ZooKeeper data must be protected using disk encryption to ensure that stored keys cannot be retrieved even if the store is compromised.”
► “In deployments without a secure ZooKeeper, SCRAM callbacks can be used to integrate with a secure external credential store.”
The dependency to remember: built-in SCRAM's security is bounded by ZooKeeper's security. That's the price of not running a password server.
3.6 SASL/OAUTHBEARER
"RFC-7628 defines the OAUTHBEARER SASL mechanism that enables credentials obtained using OAuth 2.0 to access protected resources in non-HTTP protocols. OAUTHBEARER avoids security vulnerabilities in mechanisms that use LONG-TERM PASSWORDS by using OAuth 2.0 bearer tokens with a SHORTER LIFETIME and LIMITED RESOURCE ACCESS."
⚠️ "The built-in implementation of OAUTHBEARER uses unsecured JSON Web Tokens (JWTs) and IS NOT SUITABLE FOR PRODUCTION USE. Custom callbacks can be added to integrate with standard OAuth servers."
Built-in (dev only) — note it does not validate tokens:
# BROKER
sasl.enabled.mechanisms=OAUTHBEARER
sasl.mechanism.inter.broker.protocol=OAUTHBEARER
listener.name.external.oauthbearer.sasl.jaas.config=\
org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule \
required unsecuredLoginStringClaim_sub="kafka";# CLIENT — User:Alice is the resulting default KafkaPrincipal
sasl.mechanism=OAUTHBEARER
sasl.jaas.config=\
org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule \
required unsecuredLoginStringClaim_sub="Alice";"The option
unsecuredLoginStringClaim_subis the SUBJECT CLAIM that determines theKafkaPrincipalfor the connection by default."
Production — two callbacks required:
① Client sasl.login.callback.handler.class — "to acquire tokens from the OAuth server using the long-term password or a refresh token":
@Override
public void handle(Callback[] callbacks) throws UnsupportedCallbackException {
OAuthBearerToken token = null;
for (Callback callback : callbacks) {
if (callback instanceof OAuthBearerTokenCallback) {
token = acquireToken(); // ① from OAuth server
((OAuthBearerTokenCallback) callback).token(token);
} else if (callback instanceof SaslExtensionsCallback) {// ② optional
((SaslExtensionsCallback) callback).extensions(processExtensions(token));
} else
throw new UnsupportedCallbackException(callback);
}
}⚠️ "If OAUTHBEARER is used for inter-broker communication, BROKERS must ALSO be configured with a LOGIN callback handler to acquire tokens for client connections created by the broker."
② Broker listener.name.<name>.oauthbearer.sasl.server.callback.handler.class — for validating tokens:
@Override
public void handle(Callback[] callbacks) throws UnsupportedCallbackException {
for (Callback callback : callbacks) {
if (callback instanceof OAuthBearerValidatorCallback) {
OAuthBearerValidatorCallback cb = (OAuthBearerValidatorCallback) callback;
try {
cb.token(validatedToken(cb.tokenValue())); // ① validate
} catch (OAuthBearerIllegalTokenException e) {
OAuthBearerValidationResult r = e.reason();
cb.error(errorStatus(r), r.failureScope(), r.failureOpenIdConfig());
}
} else if (callback instanceof OAuthBearerExtensionsValidatorCallback) {
OAuthBearerExtensionsValidatorCallback ecb =
(OAuthBearerExtensionsValidatorCallback) callback;
ecb.inputExtensions().map().forEach((k, v) ->
ecb.valid(validateExtension(k, v))); // ② validate extensions
} else { throw new UnsupportedCallbackException(callback); }
}
}OAUTHBEARER security considerations
- ⚠ “bearer tokens… may be used to impersonate clients → TLS must be enabled to encrypt authentication traffic.”
- ✔ “Short-lived tokens can be used to limit exposure if tokens are compromised.”
- ✔ “Reauthentication may be enabled for brokers to prevent connections outliving the tokens used for authentication. A reauthentication interval configured on brokers, combined with token revocation support, limit the amount of time an existing connection may continue to use a token after revocation.”
3.7 Delegation tokens
"Delegation tokens are shared secrets between Kafka brokers and clients that provide a lightweight configuration mechanism WITHOUT THE REQUIREMENT TO DISTRIBUTE SSL KEY STORES OR KERBEROS KEYTABS to client applications. Delegation tokens can be used to REDUCE THE LOAD ON AUTHENTICATION SERVERS, like the Kerberos Key Distribution Center (KDC)."
- The use case: “Frameworks like Kafka Connect can use delegation tokens to simplify security configuration for workers. A client that has authenticated with Kafka brokers can create delegation tokens for the same user principal and distribute these tokens to workers, which can then authenticate directly with Kafka brokers.”
- The mechanism: each token is a token identifier plus an HMAC (a shared secret). Authentication then uses SASL/SCRAM with
username = token identifierandpassword = HMAC.
# create — "If Alice runs this command, the generated token can be used to
# IMPERSONATE ALICE. The owner of this token is User:Alice."
bin/kafka-delegation-tokens.sh --bootstrap-server localhost:9092 \
--command-config admin.props --create --max-life-time-period -1 \
--renewer-principal User:Bob
# renew — "can be run by the token OWNER (Alice) or the token RENEWER (Bob)"
bin/kafka-delegation-tokens.sh --bootstrap-server localhost:9092 \
--command-config admin.props --renew --renew-time-period -1 --hmac c2VjcmV0⚠️ "To create delegation tokens for the principal
User:Alice, the client must be authenticated using Alice's credentials for any authentication protocol OTHER THAN delegation tokens. CLIENTS AUTHENTICATED USING DELEGATION TOKENS CANNOT CREATE OTHER DELEGATION TOKENS." (No privilege chaining.)
Configuration:
# ALL brokers must share the same master key
delegation.token.master.key=<key>⚠️ "This key can only be rotated by RESTARTING ALL BROKERS. ALL EXISTING TOKENS SHOULD BE DELETED BEFORE UPDATING THE MASTER KEY since they can no longer be used, and new tokens should be created AFTER the key is updated on all brokers."
"At least one of the SASL/SCRAM mechanisms must be enabled on brokers to support authentication using delegation tokens."
# CLIENT — note tokenauth="true"
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule \
required tokenauth="true" username="MTIz" password="c2VjcmV0";"The
KafkaPrincipalfor connections using this configuration will be THE ORIGINAL PRINCIPAL associated with the token, e.g.,User:Alice."
Security considerations: "suitable for production use only in deployments where ZooKeeper is secure. All the security considerations described under SCRAM also apply. The MASTER KEY used by brokers for generating tokens must be protected using encryption or by externalizing the key in a secure password store. SHORT-LIVED delegation tokens can be used to limit exposure. Reauthentication can be enabled to prevent connections operating with expired tokens."
3.8 ⚠️ Reauthentication — and the compromised-user problem
The default behavior is the problem:
"Kafka brokers perform client authentication when a connection is established. ... Kafka uses a background login thread to acquire new credentials before the old ones expire, but the new credentials are used ONLY TO AUTHENTICATE NEW CONNECTIONS by default. EXISTING CONNECTIONS THAT WERE AUTHENTICATED WITH OLD CREDENTIALS CONTINUE TO PROCESS REQUESTS until disconnection occurs due to a request timeout, an idle timeout, or network errors. LONG-LIVED CONNECTIONS MAY CONTINUE TO PROCESS REQUESTS LONG AFTER THE CREDENTIALS USED TO AUTHENTICATE THE CONNECTIONS EXPIRE."
The fix:
connections.max.reauth.ms=<positive integer>*"When set to a positive integer, Kafka brokers determine the SESSION LIFETIME for SASL connections and inform clients of this lifetime DURING THE SASL HANDSHAKE.
Session lifetime = the LOWER of (remaining lifetime of the credential) and (
connections.max.reauth.ms).Any connection that doesn't reauthenticate within this interval IS TERMINATED BY THE BROKER."*
Four scenarios it improves:
| Scenario | What reauthentication gives you |
|---|---|
| ① GSSAPI / OAUTHBEARER — limited-lifetime credentials | “guarantees that all active connections are associated with valid credentials. Short-lived credentials limit exposure.” |
| ② PLAIN / SCRAM password rotation | “Reauthentication limits the amount of time requests are processed on connections authenticated with the old password. Custom server callback that allows both old and new passwords for a period of time can be used to avoid outages until all clients migrate.” |
| ③ Nonexpiring credentials | “connections.max.reauth.ms forces reauthentication in all SASL mechanisms, including those with nonexpiring credentials. This limits the amount of time a credential may be associated with an active connection after it has been revoked.” |
| ④ Clients without reauthentication support | “terminated on session expiry, forcing the clients to reconnect and authenticate again, thus providing the same security guarantees.” |
⚠️ COMPROMISED USERS — the incident-response procedure
*"If a user is compromised, action must be taken to remove the user from the system AS SOON AS POSSIBLE. All new connections will fail to authenticate once the user is removed from the authentication server. EXISTING connections will continue to process requests until the next reauthentication timeout. If
connections.max.reauth.msIS NOT CONFIGURED, NO TIMEOUT IS APPLIED AND EXISTING CONNECTIONS MAY CONTINUE TO USE THE COMPROMISED USER'S IDENTITY FOR A LONG TIME.Kafka does not support SSL renegotiation due to known vulnerabilities... Newer protocols like TLSv1.3 do not support renegotiation. SO, EXISTING SSL CONNECTIONS MAY CONTINUE TO USE REVOKED OR EXPIRED CERTIFICATES.
💡
DenyACLs for the user principal can be used to PREVENT THESE CONNECTIONS FROM PERFORMING ANY OPERATION. SINCE ACL CHANGES ARE APPLIED WITH VERY SMALL LATENCIES ACROSS ALL BROKERS, THIS IS THE QUICKEST WAY TO DISABLE ACCESS FOR COMPROMISED USERS."*
Incident playbook — credential compromised:
- Add a Deny ACL for the principal — fastest; works on live connections, including SSL.
- Remove the user from the auth server — stops new connections.
- Wait out
connections.max.reauth.ms— kills existing SASL connections. If unset, SSL/SASL connections may persist indefinitely.
11.2 Security protocols — the 2×2 matrix
That last point matters operationally: a client bootstrapping against the EXTERNAL port will only ever be told about EXTERNAL endpoints.
11.4 Security updates without downtime
Note the shape: enable both → migrate clients → migrate inter-broker → remove the old.