如何用MLRun的KafkaSource重写带SSL的Confluent Kafka连接代码?
用MLRun KafkaSource连接Confluent Kafka的实现方案
你之前通过kafka-python实现了Confluent Kafka的连接,现在可以通过MLRun的KafkaSource类完成同样的连接逻辑,以下是对应重写后的内容:
原kafka-python连接代码
# code with usage 'kafka-python>=2.0.2' from kafka import KafkaProducer, KafkaConsumer consumer = KafkaConsumer( 'ak47-data.v1', bootstrap_servers =[ 'cpkafka01.eu.prod:9092', 'cpkafka02.eu.prod:9092', 'cpkafka03.eu.prod:9092' ], client_id='test', auto_offset_reset='earliest', sasl_mechanism="SCRAM-SHA-256", sasl_plain_password="***********", sasl_plain_username="***********", security_protocol='SASL_SSL', ssl_cafile="/v3io/bigdata/rootca.crt", ssl_certfile=None, ssl_keyfile=None) # print first topic for message in consumer: print ("%s:%d:%d: key=%s value=%s" % (message.topic, message.partition, message.offset, message.key, message.value)) break
MLRun KafkaSource重写代码
确保已安装MLRun相关依赖后,使用以下代码:
import mlrun from mlrun.datastore.sources import KafkaSource # 初始化KafkaSource,映射原代码的配置参数 kafka_source = KafkaSource( name="confluent-kafka-source", topics=["ak47-data.v1"], bootstrap_servers="cpkafka01.eu.prod:9092,cpkafka02.eu.prod:9092,cpkafka03.eu.prod:9092", auto_offset_reset="earliest", client_id="test", # SASL认证配置 sasl_mechanism="SCRAM-SHA-256", sasl_plain_username="***********", sasl_plain_password="***********", # SSL安全配置 security_protocol="SASL_SSL", ssl_cafile="/v3io/bigdata/rootca.crt" ) # 读取并处理第一条消息 stream = kafka_source.get_stream() for message in stream: print(f"{message.topic}:{message.partition}:{message.offset}: key={message.key} value={message.value}") break
参数对应说明
topics:对应原代码中KafkaConsumer指定的目标主题bootstrap_servers:原代码中服务器列表的逗号分隔字符串形式auto_offset_reset、client_id:直接匹配原代码同名参数- SASL相关参数:
sasl_mechanism、sasl_plain_username、sasl_plain_password完全沿用原配置 - SSL相关参数:
security_protocol设为SASL_SSL,ssl_cafile指定CA证书路径;原代码中ssl_certfile和ssl_keyfile为None,此处无需额外配置
内容的提问来源于stack exchange,提问作者JIST
相关产品推荐
相关产品推荐

