
Keycloak 26.1.1 OAuth2/OIDC 認証を備えた本番環境対応の Apache Kafka 4.1.0(KRaft モード)。Strimzi Kafka イメージを使用。
これは以前の POC を大幅に改善した進化版です:
Apache Kafka 4.0.0+ では、SSRF/任意ファイル読み取りの脆弱性を修正するため、JVM システムプロパティとして URL 許可リスト(org.apache.kafka.sasl.oauthbearer.allowed.urls)が導入されました。これにより、ネイティブ 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 レルムとクライアントのセットアップ
./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/ブローカー間)
↔ kafka-broker:29093 (PLAINTEXT/KRaft コントローラー)
kafka-security/ca-cert + ca-keykafka-security/broker/kafka.server.keystore.jks(サーバー証明書 + 秘密鍵を含む)kafka-security/broker/kafka.server.truststore.jks(CA 証明書を含む)changeit(すべてのキーストア/トラストストア)# ブローカー証明書
CN=kafka-broker
SAN=DNS:kafka-broker,DNS:localhost,IP:127.0.0.1
# 有効期間: 3650 日
# 鍵アルゴリズム: RSA 2048 ビット
# 署名アルゴリズム: SHA256withRSA
kafka-broker(機密)
kafka-brokersetup-keycloak.sh によって自動生成aud クレームに kafka-broker を追加preferred_username を含めるkafka-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"
}
# ノード ID
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: トークン取得用のクライアント識別子oauth.client.secret: トークン取得用のクライアントシークレットoauth.token.endpoint.uri: Keycloak トークンエンドポイント(ブローカーは内部ホスト名 keycloak:8080 を使用)oauth.valid.issuer.uri: 期待される JWT issuer(トークンの 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 は、sasl.oauthbearer.method=oidc で OAuth を実装する librdkafka(C ライブラリ)を使用します。この実装は、ネイティブ 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
期待される出力:
tcp6 0.0.0.0:9093 LISTEN (SASL_SSL)
tcp6 0.0.0.0:19092 LISTEN (PLAINTEXT)
tcp6 0.0.0.0:29093 LISTEN (CONTROLLER)
docker exec kafka-broker cat /var/lib/kafka/data/meta.properties
期待される出力:
version=1
cluster.id=kafka-cluster-01
node.id=1
問題: {"status":"invalid_token"}
oauth.jwks.endpoint.uri がブローカーコンテナから到達可能か確認docker exec kafka-broker curl http://keycloak:8080/realms/kafka-realm/protocol/openid-connect/certs問題: Token audience mismatch
aud クレームに kafka-broker が含まれていない./scripts/setup-keycloak.sh を実行して audience マッパーを追加aud クレームに kafka-broker が含まれているか確認問題: Token issuer mismatch
iss が oauth.valid.issuer.uri と一致しないoauth.valid.issuer.uri=http://localhost:8080/realms/kafka-realm(外部ホスト名)を確認http://keycloak:8080 を使用するが、issuer は http://localhost:8080 に対して検証する問題: ネイティブ Java Kafka クライアントが URL 許可リストエラーで失敗
Keycloak の JWT トークンは 5 分で期限切れになります。Strimzi OAuth は更新を自動処理します:
oauth.refresh.token: 未使用(client_credentials グラント)sasl.oauthbearer.jwks.endpoint.refresh.ms=3600000 # 1 時間
sasl.oauthbearer.jwks.endpoint.retry.backoff.ms=100
sasl.oauthbearer.jwks.endpoint.retry.backoff.max.ms=10000
connections.max.idle.ms=600000
connection.failed.authentication.delay.ms=1000
ssl.endpoint.identification.algorithm=https に更新(none を削除)allow.everyone.if.no.acl.found=true を削除)kafka-acls --bootstrap-server localhost:9093 \
--command-config admin.properties \
--add --allow-principal User:kafka-producer \
--operation Write --topic '*'
oauth.token.endpoint.uri と oauth.jwks.endpoint.uri を HTTPS URL に更新.
├── docker-compose.yml # オーケストレーション
├── .env # シークレット(gitignore 対象)
├── kafka-config/
│ ├── kraft-config.properties # Kafka ブローカー設定
│ ├── producer.properties # プロデューサー OAuth 設定(CLI ツール用)
│ └── consumer.properties # コンシューマー OAuth 設定(CLI ツール用)
├── kafka-security/
│ ├── generate-certs.sh # SSL 証明書生成スクリプト
│ ├── ca-cert # ルート CA 証明書
│ ├── ca-key # ルート CA 秘密鍵
│ └── broker/
│ ├── kafka.server.keystore.jks
│ └── kafka.server.truststore.jks
├── scripts/
│ └── setup-keycloak.sh # Keycloak レルム/クライアントセットアップ
└── tests/
└── quick_test.py # OAuth 検証テスト
Strimzi Kafka イメージ(quay.io/strimzi/kafka:0.48.0-kafka-4.1.0)を公式 Apache Kafka イメージの代わりに使用する理由:
io.strimzi.kafka.oauth.*)イメージ構成:
ブローカー設定には 2 つの URL があります:
oauth.token.endpoint.uri=http://keycloak:8080/...(内部 Docker ネットワーク)oauth.valid.issuer.uri=http://localhost:8080/...(外部、JWT の iss クレームと一致)これは以下の理由によるものです:
ブローカーは JWT の preferred_username クレームからプリンシパルを抽出します:
service-account-kafka-producer → User:service-account-kafka-producer
ACL はこのプリンシパルを参照して認可を行います。
| コンポーネント | バージョン | 備考 |
|---|
| Apache Kafka | 4.1.0 | KRaft モード(ZooKeeper なし) |
| Strimzi Kafka イメージ | 0.48.0 | Docker イメージ: quay.io/strimzi/kafka:0.48.0-kafka-4.1.0 |
| Strimzi OAuth ライブラリ | 0.17.0 | Strimzi Kafka 0.48.0 イメージにプリインストール |
| Keycloak | 26.1.1 | 最新 LTS |
| librdkafka | 2.12.0+ | OIDC OAuth サポート |
| confluent-kafka-python | 2.12.0+ | librdkafka バージョンと一致 |