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

在ODI 12C中以Spark为中转从Kafka映射到CSV文件时遇错误

解决ODI 12C中Spark中转Kafka到CSV时的AttributeError问题

这个错误我之前帮人排查过,核心是PySpark Streaming Kafka的参数传错了,咱们一步步来理清楚:

错误原因拆解

从报错里的代码片段能看到关键问题:

,"group.id" : "kafka_file_Physical"}, fromOffsets='earliest')

这里的fromOffsets='earliest'完全用错了地方:

  • PySpark的createDirectStream方法要求fromOffsets是字典类型,格式是{主题名: 分区偏移量},用来指定特定分区的起始消费位置;
  • 你传了字符串'earliest',而字符串根本没有items()方法(Spark内部会调用这个方法解析偏移量配置),所以直接抛出了AttributeError: 'str' object has no attribute 'items'。

如果你的需求是从Kafka的最早可用偏移量开始消费,这个配置应该放到kafkaParams字典里,而不是传给fromOffsets。

修复方案

方案1:从最早偏移量开始消费(最常见需求)

修改PySpark代码的参数配置,把auto.offset.reset加到kafka参数中,移除错误的fromOffsets='earliest':

# 正确的Kafka参数配置
kafkaParams = {
    "group.id": "kafka_file_Physical",
    "auto.offset.reset": "earliest"  # 这里配置从最早偏移量启动消费
}
# 调用createDirectStream时无需传fromOffsets
stream = KafkaUtils.createDirectStream(ssc, ["你的主题名"], kafkaParams)

方案2:指定具体分区的偏移量(如果需要精准控制)

如果你确实需要从某个分区的特定位置开始消费,fromOffsets必须传字典,示例如下:

from pyspark.streaming.kafka import TopicAndPartition

# 定义特定主题分区的起始偏移量
fromOffsets = {
    TopicAndPartition("你的主题名", 0): 100,  # 0号分区从第100条开始
    TopicAndPartition("你的主题名", 1): 200   # 1号分区从第200条开始
}
# 调用时传入正确的字典格式参数
stream = KafkaUtils.createDirectStream(ssc, ["你的主题名"], kafkaParams, fromOffsets=fromOffsets)

最后回到ODI 12C里,找到生成这个PySpark脚本的Spark EKM配置节点,调整Kafka消费的偏移量配置方式,确保参数传递符合PySpark的要求就行。


内容的提问来源于stack exchange,提问作者rounak chourasia

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 06:55:57