生产就绪的 Apache Kafka 4.1.0(KRaft 模式),使用 Strimzi Kafka 镜像实现 Keycloak 26.1.1 OAuth2/OIDC 认证。
这是之前 POC 的演进版本,具有显著改进:
Apache Kafka 4.0.0+ 引入了 URL 白名单(org.apache.kafka.sasl.oauthbearer.allowed.urls)作为 JVM 系统属性,以修复 SSRF/任意文件读取漏洞。这破坏了原生 Apache Kafka 客户端中的标准 OAuth 用法。
解决方案:Strimzi Kafka OAuth 库未实现此限制,从而在 Kafka 4.1.0 上启用 OAuth 功能。
# 生成 SSL 证书
cd kafka-security
./generate-certs.sh
cd ..
# 启动服务
docker compose up -d
# 验证 Keycloak
curl http://localhost:8080/health/ready
# 设置 Keycloak realm 和客户端
./scripts/setup-keycloak.sh
# 测试 OAuth 生产者
source ~/.venv/bin/activate
uv pip install confluent-kafka
python tests/quick_test.py
keycloak:8080 (HTTP) ←→ kafka-broker:9093 (SASL_SSL/OAuth)
↔ kafka-broker:19092 (PLAINTEXT/inter-broker)
↔ kafka-broker:29093 (PLAINTEXT/KRaft controller)
kafka-security/ca-cert + ca-keykafka-security/broker/kafka.server.keystore.jks(包含服务器证书 + 私钥)kafka-security/broker/kafka.server.truststore.jks(包含 CA 证书)changeit(所有密钥库/信任库)# Broker 证书
CN=kafka-broker
SAN=DNS:kafka-broker,DNS:localhost,IP:127.0.0.1
# 有效期:3650 天
# 密钥算法:RSA 2048 位
# 签名算法:SHA256withRSA
kafka-broker(机密)
kafka-brokersetup-keycloak.sh 自动生成kafka-broker 添加到 JWT aud 声明preferred_usernamekafka-producer(机密)
kafka-producerclient_credentialskafka-consumer(机密)
kafka-consumerclient_credentialsPOST http://localhost:8080/realms/kafka-realm/protocol/openid-connect/token
Content-Type: application/x-www-form-urlencoded
grant_type=client_credentials
&client_id=kafka-producer
&client_secret=<secret>
&scope=profile email
{
"aud": ["kafka-broker", "account"],
"iss": "http://localhost:8080/realms/kafka-realm",
"azp": "kafka-producer",
"preferred_username": "service-account-kafka-producer",
"scope": "profile email"
}
# 节点标识
node.id=1
process.roles=broker,controller
controller.quorum.voters=1@kafka-broker:29093
# 监听器
listeners=SASL_SSL://0.0.0.0:9093,PLAINTEXT://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093
advertised.listeners=SASL_SSL://localhost:9093,PLAINTEXT://kafka-broker:19092
listener.security.protocol.map=SASL_SSL:SASL_SSL,PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
inter.broker.listener.name=PLAINTEXT
controller.listener.names=CONTROLLER
# SASL 机制
sasl.enabled.mechanisms=OAUTHBEARER
# Strimzi OAuth 处理器(针对 SASL_SSL 的按监听器配置)
listener.name.sasl_ssl.oauthbearer.sasl.login.callback.handler.class=io.strimzi.kafka.oauth.client.JaasClientOauthLoginCallbackHandler
listener.name.sasl_ssl.oauthbearer.sasl.server.callback.handler.class=io.strimzi.kafka.oauth.server.JaasServerOauthValidatorCallbackHandler
# 通过 JAAS 的 OAuth 配置
listener.name.sasl_ssl.oauthbearer.sasl.jaas.config=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \
oauth.client.id="kafka-broker" \
oauth.client.secret="<secret>" \
oauth.token.endpoint.uri="http://keycloak:8080/realms/kafka-realm/protocol/openid-connect/token" \
oauth.valid.issuer.uri="http://localhost:8080/realms/kafka-realm" \
oauth.jwks.endpoint.uri="http://keycloak:8080/realms/kafka-realm/protocol/openid-connect/certs" \
oauth.username.claim="preferred_username";
oauth.client.id:用于获取 token 的客户端标识符oauth.client.secret:用于获取 token 的客户端密钥oauth.token.endpoint.uri:Keycloak token 端点(broker 使用内部主机名 keycloak:8080)oauth.valid.issuer.uri:预期的 JWT issuer(必须与 token iss 声明匹配,使用外部 localhost:8080)oauth.jwks.endpoint.uri:用于 JWT 签名验证的 JWKS 端点oauth.username.claim:用于主体提取的 JWT 声明authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
super.users=User:kafka-broker;User:ANONYMOUS
allow.everyone.if.no.acl.found=true
注意:目前为测试目的采用宽松配置。生产环境应使用 ACL。
from confluent_kafka import Producer
conf = {
'bootstrap.servers': 'localhost:9093',
'security.protocol': 'SASL_SSL',
'sasl.mechanisms': 'OAUTHBEARER',
'sasl.oauthbearer.method': 'oidc',
'sasl.oauthbearer.client.id': 'kafka-producer',
'sasl.oauthbearer.client.secret': '<secret>',
'sasl.oauthbearer.token.endpoint.url': 'http://localhost:8080/realms/kafka-realm/protocol/openid-connect/token',
'ssl.ca.location': 'kafka-security/ca-cert',
'ssl.endpoint.identification.algorithm': 'none',
}
producer = Producer(conf)
producer.produce('topic', b'message')
producer.flush()
from confluent_kafka import Consumer
conf = {
'bootstrap.servers': 'localhost:9093',
'group.id': 'test-group',
'security.protocol': 'SASL_SSL',
'sasl.mechanisms': 'OAUTHBEARER',
'sasl.oauthbearer.method': 'oidc',
'sasl.oauthbearer.client.id': 'kafka-consumer',
'sasl.oauthbearer.client.secret': '<secret>',
'sasl.oauthbearer.token.endpoint.url': 'http://localhost:8080/realms/kafka-realm/protocol/openid-connect/token',
'ssl.ca.location': 'kafka-security/ca-cert',
'ssl.endpoint.identification.algorithm': 'none',
'auto.offset.reset': 'earliest',
}
consumer = Consumer(conf)
consumer.subscribe(['topic'])
while True:
msg = consumer.poll(1.0)
if msg: print(msg.value())
confluent-kafka-python 使用 librdkafka(C 库),通过 sasl.oauthbearer.method=oidc 实现 OAuth。该实现不检查阻止原生 Apache Kafka Java 客户端的 org.apache.kafka.sasl.oauthbearer.allowed.urls 系统属性。
TOKEN=$(curl -s -X POST http://localhost:8080/realms/kafka-realm/protocol/openid-connect/token \
-d "grant_type=client_credentials" \
-d "client_id=kafka-producer" \
-d "client_secret=<secret>" | jq -r .access_token)
echo $TOKEN | cut -d. -f2 | base64 -d 2>/dev/null | jq .
预期声明:
{
"aud": ["kafka-broker", "account"],
"iss": "http://localhost:8080/realms/kafka-realm",
"azp": "kafka-producer",
"preferred_username": "service-account-kafka-producer"
}
docker logs kafka-broker 2>&1 | grep -E "Strimzi|JWTSignatureValidator|OAUTHBEARER"
预期输出:
[io.strimzi.kafka.oauth.validator.JWTSignatureValidator] JWKS keys change detected
docker exec kafka-broker netstat -tlnp | grep java