You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 04:45:31