Confluent-Kafka Python客户端Kerberos认证问题(Cloudera集群)
我之前帮同事解决过类似的问题,给你整理一套适配Cloudera Kafka+Kerberos的confluent-kafka客户端配置步骤,亲测有效:
配置confluent-kafka客户端对接Kerberos认证的Cloudera Kafka
一、先搞定本地Kerberos基础环境
- 首先在你的客户端机器上安装Kerberos客户端工具,比如Debian/Ubuntu用
apt install krb5-user,RHEL/CentOS用yum install krb5-workstation,这是认证的前提。 - 关键的
krb5.conf配置直接从Cloudera集群节点上拷贝过来用!手动配置很容易出错,集群的krb5.conf一般在/etc/krb5.conf,拷贝到你的客户端机器同路径下就行。如果非要手动写,核心要包含:[realms] # 替换成你的集群Kerberos Realm YOUR_REALM.COM = { kdc = 你的KDC服务器地址 admin_server = 你的KDC服务器地址 default_domain = your-realm.com } [domain_realm] .your-realm.com = YOUR_REALM.COM your-realm.com = YOUR_REALM.COM
二、获取有效的Kerberos Ticket
- 临时测试的话,直接用密码获取Ticket:
输入密码后,用kinit your-kerberos-username@YOUR_REALM.COMklist命令确认Ticket是否生效,能看到Valid starting和Expires时间就没问题。 - 如果是长期运行的服务,建议用keytab文件(找集群管理员给你生成),避免每次手动输密码:
kinit -kt /path/to/your/user.keytab your-kerberos-username@YOUR_REALM.COM
三、配置confluent-kafka客户端参数
不管是生产者还是消费者,都要在配置里明确Kerberos相关参数,以下是实际可用的示例:
生产者配置(Python)
from confluent_kafka import Producer # 替换成你的集群Broker地址、Kerberos信息 conf = { 'bootstrap.servers': 'kafka-broker1.your-realm.com:9092,kafka-broker2.your-realm.com:9092', 'security.protocol': 'SASL_PLAINTEXT', # 集群如果用SSL加密就改成SASL_SSL 'sasl.mechanism': 'GSSAPI', 'sasl.kerberos.service.name': 'kafka', # Cloudera Kafka默认服务名是kafka,别改! # 用keytab的话加上下面两行,不用的话可以注释掉 'sasl.kerberos.keytab': '/path/to/your/user.keytab', 'sasl.kerberos.principal': 'your-kerberos-username@YOUR_REALM.COM' } producer = Producer(conf) # 测试生产一条消息 def delivery_report(err, msg): if err is not None: print(f'Message delivery failed: {err}') else: print(f'Message delivered to {msg.topic()} [{msg.partition()}]') producer.produce('test-topic', key='key1', value='test-value', callback=delivery_report) producer.flush()
消费者配置(Python)
from confluent_kafka import Consumer, KafkaError conf = { 'bootstrap.servers': 'kafka-broker1.your-realm.com:9092,kafka-broker2.your-realm.com:9092', 'group.id': 'test-consumer-group', 'auto.offset.reset': 'earliest', 'security.protocol': 'SASL_PLAINTEXT', 'sasl.mechanism': 'GSSAPI', 'sasl.kerberos.service.name': 'kafka', 'sasl.kerberos.keytab': '/path/to/your/user.keytab', 'sasl.kerberos.principal': 'your-kerberos-username@YOUR_REALM.COM' } consumer = Consumer(conf) consumer.subscribe(['test-topic']) while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: print(f'Reached end of partition {msg.partition()}') else: print(f'Consumer error: {msg.error()}') break else: print(f'Received message: {msg.value().decode("utf-8")}') consumer.close()
四、常见坑点排查
- 如果报认证失败,先查
klist看Ticket有没有过期,过期就重新kinit。 sasl.kerberos.service.name一定要是kafka,Cloudera集群默认配置就是这个,写错了肯定连不上。- 检查客户端机器能不能访问KDC服务器和Kafka Broker的端口(默认9092,SSL的话是9093),防火墙别挡住了。
- 如果是SASL_SSL的场景,还要加上SSL配置:
'ssl.ca.location': '/path/to/ca-cert.pem',证书找集群管理员要。
内容的提问来源于stack exchange,提问作者Dave
相关产品推荐
相关产品推荐

